diff --git a/cachelib/CMakeLists.txt b/cachelib/CMakeLists.txt index 29f3ee2c7f..ef545f1c12 100644 --- a/cachelib/CMakeLists.txt +++ b/cachelib/CMakeLists.txt @@ -44,6 +44,9 @@ set(CMAKE_POSITION_INDEPENDENT_CODE ON) option(BUILD_TESTS "If enabled, compile the tests." ON) +option(BUILD_WITH_DTO + "If enabled, build with the DTO library for Intel DSA offload support." OFF) + set(BIN_INSTALL_DIR bin CACHE STRING "The subdirectory where binaries should be installed") @@ -312,6 +315,56 @@ endfunction() add_thrift_file(OBJECT_CACHE_PERSISTENCE object_cache/persistence/persistent_data.thrift json) +if (BUILD_WITH_DTO) + # DSA offload runtime (DTO). Resolved here so every cachelib library can + # use it: DTO::dto is linked PUBLIC into cachelib_common (see + # common/CMakeLists.txt), the base library of every other cachelib target, + # which propagates the headers, the link dependency, and the + # CACHELIB_BUILD_WITH_DTO compile definition to the whole tree and to + # consumers of the installed package. + # + # The DSA checksum offload code requires a DTO with the caller-allocated + # async submit/poll API and the raw-CRC32C convention, which upstream + # intel/DTO does not have — so when no installed DTO cmake package is + # found, fetch and build the pinned revision in-tree rather than linking + # whatever libdto happens to be on the system. + set(CACHELIB_DTO_GIT_REPOSITORY "https://github.com/byrnedj/DTO.git" + CACHE STRING "Git repository for DTO when no installed copy is found") + set(CACHELIB_DTO_GIT_TAG "991f4f9084095af78e2d7ca2dcb9a8c6069f8904" + CACHE STRING "DTO revision to fetch when no installed copy is found") + + find_package(DTO CONFIG QUIET) + if (DTO_FOUND) + message(STATUS "Using installed DTO package: ${DTO_DIR}") + else() + if (CMAKE_VERSION VERSION_LESS 3.14) + message(FATAL_ERROR + "BUILD_WITH_DTO needs an installed DTO cmake package or " + "CMake >= 3.14 to fetch one. Install DTO from " + "${CACHELIB_DTO_GIT_REPOSITORY} at ${CACHELIB_DTO_GIT_TAG}, " + "or upgrade CMake.") + endif() + # DTO links against libaccel-config and libnuma; check for them here so + # a missing system package fails at configure time with a clear message + # instead of at link time inside the fetched project. + find_library(ACCEL_CONFIG_LIBRARY accel-config) + find_library(NUMA_LIBRARY numa) + if (NOT ACCEL_CONFIG_LIBRARY OR NOT NUMA_LIBRARY) + message(FATAL_ERROR + "BUILD_WITH_DTO requires the libaccel-config and libnuma " + "development packages (e.g. apt install libaccel-config-dev " + "libnuma-dev)") + endif() + message(STATUS "DTO not installed; fetching " + "${CACHELIB_DTO_GIT_REPOSITORY} @ ${CACHELIB_DTO_GIT_TAG}") + include(FetchContent) + FetchContent_Declare(dto + GIT_REPOSITORY "${CACHELIB_DTO_GIT_REPOSITORY}" + GIT_TAG "${CACHELIB_DTO_GIT_TAG}") + FetchContent_MakeAvailable(dto) + endif() +endif() + add_subdirectory (common) add_subdirectory (shm) add_subdirectory (navy) diff --git a/cachelib/allocator/Cache.cpp b/cachelib/allocator/Cache.cpp index 6ef8777f21..dd1dc79770 100644 --- a/cachelib/allocator/Cache.cpp +++ b/cachelib/allocator/Cache.cpp @@ -426,6 +426,12 @@ void CacheBase::updateGlobalCacheStats(const std::string& statPrefix) const { if (stats.nvmCacheEnabled) { counters_.updateDelta(statPrefix + "nvm.alloc_attempts", stats.numNvmAllocAttempts); + counters_.updateDelta(statPrefix + "nvm.copy_out_offloaded", + stats.numNvmCopyOutOffloaded); + counters_.updateDelta(statPrefix + "nvm.copy_out_offloaded_bytes", + stats.numNvmCopyOutOffloadedBytes); + counters_.updateDelta(statPrefix + "nvm.copy_out_fallbacks", + stats.numNvmCopyOutFallbacks); counters_.updateDelta(statPrefix + "nvm.destructor_alloc", stats.numNvmAllocForItemDestructor); counters_.updateDelta(statPrefix + "nvm.destructor_alloc_errors", diff --git a/cachelib/allocator/CacheStats.cpp b/cachelib/allocator/CacheStats.cpp index f961a49a0f..a8f95d5fe0 100644 --- a/cachelib/allocator/CacheStats.cpp +++ b/cachelib/allocator/CacheStats.cpp @@ -54,10 +54,10 @@ struct SizeVerify {}; void Stats::populateGlobalCacheStats(GlobalCacheStats& ret) const { #ifndef SKIP_SIZE_VERIFY #ifdef __GLIBCXX__ -#define EXPECTED_SIZE 16944 +#define EXPECTED_SIZE 17200 #endif #ifdef _LIBCPP_VERSION -#define EXPECTED_SIZE 16944 +#define EXPECTED_SIZE 17200 #endif SizeVerify a = SizeVerify{}; std::ignore = a; @@ -106,6 +106,9 @@ void Stats::populateGlobalCacheStats(GlobalCacheStats& ret) const { ret.numInsertOrReplaceInserted = numInsertOrReplaceInserted.get(); ret.numInsertOrReplaceReplaced = numInsertOrReplaceReplaced.get(); ret.numNvmAllocAttempts = numNvmAllocAttempts.get(); + ret.numNvmCopyOutOffloaded = numNvmCopyOutOffloaded.get(); + ret.numNvmCopyOutOffloadedBytes = numNvmCopyOutOffloadedBytes.get(); + ret.numNvmCopyOutFallbacks = numNvmCopyOutFallbacks.get(); ret.numNvmAllocForItemDestructor = numNvmAllocForItemDestructor.get(); ret.numNvmItemDestructorAllocErrors = numNvmItemDestructorAllocErrors.get(); diff --git a/cachelib/allocator/CacheStats.h b/cachelib/allocator/CacheStats.h index 6187e8d51b..d5adc779d1 100644 --- a/cachelib/allocator/CacheStats.h +++ b/cachelib/allocator/CacheStats.h @@ -456,6 +456,11 @@ struct GlobalCacheStats { // attempts made from nvm cache to allocate an item for promotion uint64_t numNvmAllocAttempts{0}; + // flash-hit copy-outs performed by Intel DSA, bytes moved, CPU fallbacks + uint64_t numNvmCopyOutOffloaded{0}; + uint64_t numNvmCopyOutOffloadedBytes{0}; + uint64_t numNvmCopyOutFallbacks{0}; + // attempts made from nvm cache to allocate an item for its destructor uint64_t numNvmAllocForItemDestructor{0}; // heap allocate errors for item destructor diff --git a/cachelib/allocator/CacheStatsInternal.h b/cachelib/allocator/CacheStatsInternal.h index b3225de37b..53573d2dcf 100644 --- a/cachelib/allocator/CacheStatsInternal.h +++ b/cachelib/allocator/CacheStatsInternal.h @@ -141,6 +141,12 @@ struct Stats { // attempts made from nvm cache to allocate an item for promotion TLCounter numNvmAllocAttempts{0}; + // flash-hit copy-outs (Navy buffer -> DRAM item) performed by Intel DSA, + // the bytes they moved, and offload attempts that fell back to the CPU + TLCounter numNvmCopyOutOffloaded{0}; + TLCounter numNvmCopyOutOffloadedBytes{0}; + TLCounter numNvmCopyOutFallbacks{0}; + // attempts made from nvm cache to allocate an item for its destructor TLCounter numNvmAllocForItemDestructor{0}; // heap allocate errors for item destructor diff --git a/cachelib/allocator/nvmcache/NavyConfig.cpp b/cachelib/allocator/nvmcache/NavyConfig.cpp index 63966a4fff..d5119dfd3e 100644 --- a/cachelib/allocator/nvmcache/NavyConfig.cpp +++ b/cachelib/allocator/nvmcache/NavyConfig.cpp @@ -323,6 +323,10 @@ std::map EnginesConfig::serialize() const { folly::to(blockCache().getNumInMemBuffers()); configMap["navyConfig::blockCacheDataChecksum"] = blockCache().getDataChecksum() ? "true" : "false"; + configMap["navyConfig::blockCacheChecksumOffload"] = + blockCache().getChecksumOffload() ? "true" : "false"; + configMap["navyConfig::blockCacheChecksumOffloadMinSize"] = + folly::to(blockCache().getChecksumOffloadMinSize()); configMap["navyConfig::blockCacheSegmentedFifoSegmentRatio"] = folly::join(",", blockCache().getSFifoSegmentRatio()); @@ -335,6 +339,8 @@ std::map EnginesConfig::serialize() const { folly::to(bigHash().getBucketBfSize()); configMap["navyConfig::bigHashSmallItemMaxSize"] = folly::to(bigHash().getSmallItemMaxSize()); + configMap["navyConfig::bigHashChecksumOffload"] = + bigHash().getChecksumOffload() ? "true" : "false"; return configMap; } diff --git a/cachelib/allocator/nvmcache/NavyConfig.h b/cachelib/allocator/nvmcache/NavyConfig.h index bfe094170e..a6ed68e90b 100644 --- a/cachelib/allocator/nvmcache/NavyConfig.h +++ b/cachelib/allocator/nvmcache/NavyConfig.h @@ -536,6 +536,36 @@ class BlockCacheConfig { return *this; } + // Offload data checksumming (fused with the value copy on the write path) + // to Intel DSA via the DTO library. Requires data checksum to be enabled + // and CacheLib built with BUILD_WITH_DTO. @minSize is the minimum value + // size to use the offloaded path; smaller values are checksummed in + // software, since submitting and polling a descriptor costs about what the + // CPU needs to CRC 16-32 KiB (16 KiB measured as the break-even). + BlockCacheConfig& setChecksumOffload(bool checksumOffload, + uint32_t minSize = 16384) noexcept { + checksumOffload_ = checksumOffload; + checksumOffloadMinSize_ = minSize; + return *this; + } + + // Separate minimum size for offloading checksum VERIFICATION on the read, + // reclaim and reinsertion paths; 0 = same as the write-path gate. Verifying + // has no CPU work to overlap with the accelerator, so its break-even is + // higher than the fused write. + BlockCacheConfig& setChecksumOffloadReadMinSize(uint32_t minSize) noexcept { + checksumOffloadReadMinSize_ = minSize; + return *this; + } + + // Cache-control hint on the fused write-path copy (default true): steer the + // value bytes toward the CPU cache. Turn off when the next reader of the + // region buffer is the device (directFlush / flush copy offload). + BlockCacheConfig& setChecksumOffloadCacheControl(bool enable) noexcept { + checksumOffloadCacheControl_ = enable; + return *this; + } + BlockCacheConfig& setPreciseRemove(bool preciseRemove) noexcept { preciseRemove_ = preciseRemove; return *this; @@ -561,6 +591,17 @@ class BlockCacheConfig { return *this; } + // When directFlush is off: copy the region buffer into the flush write + // buffer with one Intel DSA Memory Move instead of memcpy. Requires + // CacheLib built with BUILD_WITH_DTO; falls back to memcpy if DSA is + // unusable. The write buffers are pooled, so after the first flush per + // buffer their pages are resident; a work queue with block-on-fault covers + // the first touch. + BlockCacheConfig& setFlushCopyOffload(bool enable) noexcept { + flushCopyOffload_ = enable; + return *this; + } + BlockCacheConfig& setAllocatorCount(uint32_t numAllocators) noexcept { allocatorsPerPriority_ = {numAllocators}; return *this; @@ -617,6 +658,18 @@ class BlockCacheConfig { bool getDataChecksum() const { return dataChecksum_; } + bool getChecksumOffload() const { return checksumOffload_; } + + uint32_t getChecksumOffloadMinSize() const { return checksumOffloadMinSize_; } + + uint32_t getChecksumOffloadReadMinSize() const { + return checksumOffloadReadMinSize_; + } + + bool getChecksumOffloadCacheControl() const { + return checksumOffloadCacheControl_; + } + uint64_t getSize() const { return size_; } bool isRegionManagerFlushAsync() const { return regionManagerFlushAsync_; } @@ -624,6 +677,7 @@ class BlockCacheConfig { bool isRecoverEvictionPolicy() const { return recoverEvictionPolicy_; } bool isDirectFlush() const { return directFlush_; } + bool isFlushCopyOffload() const { return flushCopyOffload_; } bool isCombinedEntryBlockEnabled() const { return useCombinedEntryBlock_; } @@ -660,6 +714,12 @@ class BlockCacheConfig { uint32_t regionSize_{16 * 1024 * 1024}; // Whether enabling data checksum for Navy BlockCache. bool dataChecksum_{true}; + // Whether to offload data checksumming to Intel DSA (fused with the value + // copy on the write path), and the minimum value size to do so. + bool checksumOffload_{false}; + uint32_t checksumOffloadMinSize_{16384}; + uint32_t checksumOffloadReadMinSize_{0}; + bool checksumOffloadCacheControl_{true}; // Whether to remove an item by checking the key (true) or only the hash value // (false). bool preciseRemove_{false}; @@ -677,6 +737,9 @@ class BlockCacheConfig { // Whether to write region buffer directly without intermediate copy. bool directFlush_{false}; + // Whether to do the flush copy (when not directFlush) on Intel DSA. + bool flushCopyOffload_{false}; + // Whether to use Combined entry block (For index entries and small sized // items). // Only FixedSizeIndex will support this and it doesn't work with @@ -739,6 +802,14 @@ class BigHashConfig { return *this; } + // Offload bucket checksumming to Intel DSA via the DTO library. Most + // beneficial with large buckets (16KB+). Requires CacheLib built with + // BUILD_WITH_DTO. + BigHashConfig& setChecksumOffload(bool checksumOffload) noexcept { + checksumOffload_ = checksumOffload; + return *this; + } + bool isBloomFilterEnabled() const { return bucketBfSize_ > 0; } unsigned int getSizePct() const { return sizePct_; } @@ -751,6 +822,8 @@ class BigHashConfig { uint8_t getNumMutexesPower() const { return numMutexesPower_; } + bool getChecksumOffload() const { return checksumOffload_; } + private: // Percentage of how much of the device out of all is given to BigHash // engine in Navy, e.g. 50. @@ -766,6 +839,8 @@ class BigHashConfig { uint64_t smallItemMaxSize_{}; // numMutexes = 1 << numMutexesPower_. uint8_t numMutexesPower_{14}; + // Whether to offload bucket checksumming to Intel DSA. + bool checksumOffload_{false}; }; // Config for a pair of small,large engines. diff --git a/cachelib/allocator/nvmcache/NavySetup.cpp b/cachelib/allocator/nvmcache/NavySetup.cpp index 0a184ef3d0..28f3326b4a 100644 --- a/cachelib/allocator/nvmcache/NavySetup.cpp +++ b/cachelib/allocator/nvmcache/NavySetup.cpp @@ -140,6 +140,8 @@ uint64_t setupBigHash(const navy::BigHashConfig& bigHashConfig, // Set number of mutexes from config bigHash->setNumMutexesPower(bigHashConfig.getNumMutexesPower()); + bigHash->setChecksumOffload(bigHashConfig.getChecksumOffload()); + proto.setBigHash(std::move(bigHash), bigHashConfig.getSmallItemMaxSize()); if (bigHashCacheOffset <= bigHashStartOffsetLimit) { @@ -196,6 +198,12 @@ uint64_t setupBlockCache(const navy::BlockCacheConfig& blockCacheConfig, auto blockCache = cachelib::navy::createBlockCacheProto(); blockCache->setLayout(blockCacheOffset, blockCacheSize, regionSize); blockCache->setChecksum(blockCacheConfig.getDataChecksum()); + blockCache->setChecksumOffload(blockCacheConfig.getChecksumOffload(), + blockCacheConfig.getChecksumOffloadMinSize()); + blockCache->setChecksumOffloadReadMinSize( + blockCacheConfig.getChecksumOffloadReadMinSize()); + blockCache->setChecksumOffloadCacheControl( + blockCacheConfig.getChecksumOffloadCacheControl()); // set eviction policy auto segmentRatio = blockCacheConfig.getSFifoSegmentRatio(); @@ -214,6 +222,7 @@ uint64_t setupBlockCache(const navy::BlockCacheConfig& blockCacheConfig, blockCache->setRecoverEvictionPolicy( blockCacheConfig.isRecoverEvictionPolicy()); blockCache->setDirectFlush(blockCacheConfig.isDirectFlush()); + blockCache->setFlushCopyOffload(blockCacheConfig.isFlushCopyOffload()); blockCache->setUseCombinedEntryBlock( blockCacheConfig.isCombinedEntryBlockEnabled()); blockCache->setNumAllocatorsPerPriority( diff --git a/cachelib/allocator/nvmcache/NvmCache.h b/cachelib/allocator/nvmcache/NvmCache.h index b1d7a5f11b..b7704bbdeb 100644 --- a/cachelib/allocator/nvmcache/NvmCache.h +++ b/cachelib/allocator/nvmcache/NvmCache.h @@ -34,6 +34,7 @@ #include "cachelib/allocator/nvmcache/CacheApiWrapper.h" #include "cachelib/allocator/nvmcache/InFlightPuts.h" #include "cachelib/allocator/nvmcache/NavyConfig.h" +#include "cachelib/navy/common/ChecksumOffload.h" #include "cachelib/allocator/nvmcache/NavySetup.h" #include "cachelib/allocator/nvmcache/NvmItem.h" #include "cachelib/allocator/nvmcache/ReqContexts.h" @@ -160,6 +161,17 @@ class NvmCache { // in the future. bool disableNvmCacheOnBadState{true}; + // Offload the flash-hit copy-out (Navy read buffer -> DRAM item) to Intel + // DSA through the DTO library for blobs of at least copyOutOffloadMinSize + // bytes. Requires CacheLib built with BUILD_WITH_DTO; verified at + // construction and silently falls back to memcpy when DSA is unusable. + // The DRAM cache memory should be resident (pre-touched or huge-page slabs) or + // the DSA work queue configured with block-on-fault, otherwise device + // page faults turn every offloaded copy into a failed descriptor plus a + // CPU copy. + bool copyOutOffload{false}; + uint32_t copyOutOffloadMinSize{65536}; + // serialize the config for debugging purposes std::map serialize() const; @@ -384,6 +396,11 @@ class NvmCache { // based on the NvmItem WriteHandle createItem(folly::StringPiece key, const NvmItem& nvmItem); + // Copies a blob's payload into freshly allocated item memory, on DSA when + // copy-out offload is active and the blob is large enough (see + // Config::copyOutOffload), otherwise with memcpy. + void copyBlobData(void* dest, const void* src, size_t n); + // creates the item into IOBuf from NvmItem, if the item has chained items, // chained IOBufs will be created. // @param key key for the dipper item @@ -566,6 +583,9 @@ class NvmCache { } const Config config_; + + // Whether the DSA copy-out offload passed its runtime self-check. + bool copyOutOffload_{false}; C& cache_; //< cache allocator std::atomic navyEnabled_{true}; //< switch to turn off/on navy @@ -646,6 +666,8 @@ std::map NvmCache::Config::serialize() const { truncateItemToOriginalAllocSizeInNvm ? "true" : "false"; configMap["disableNvmCacheOnBadState"] = disableNvmCacheOnBadState ? "true" : "false"; + configMap["copyOutOffload"] = copyOutOffload ? "true" : "false"; + configMap["copyOutOffloadMinSize"] = std::to_string(copyOutOffloadMinSize); return configMap; } @@ -1130,6 +1152,19 @@ NvmCache::NvmCache(C& c, itemDestructor_ ? true : false, navyPersistParams); + if (config_.copyOutOffload) { + copyOutOffload_ = navy::copyOffloadSelfCheck(); + if (copyOutOffload_) { + XLOGF(INFO, + "NvmCache: DSA copy-out offload active for blobs >= {} bytes", + config_.copyOutOffloadMinSize); + } else { + XLOG(WARN) << "NvmCache: DSA copy-out offload requested but the DTO/DSA " + "self-check failed (not built with DTO, or DSA " + "unavailable); falling back to memcpy"; + } + } + if (accessTimeMap_ && shmManager_) { try { accessTimeMap_->recover(*shmManager_); @@ -1484,7 +1519,7 @@ typename NvmCache::WriteHandle NvmCache::createItem( } else { XDCHECK_LE(pBlob.data.size(), getStorageSizeInNvm(*it)); XDCHECK_LE(pBlob.origAllocSize, pBlob.data.size()); - ::memcpy(it->getMemory(), pBlob.data.data(), pBlob.data.size()); + copyBlobData(it->getMemory(), pBlob.data.data(), pBlob.data.size()); it->markNvmClean(); if (isNvmItemLarge(key, nvmItem)) { it->markNvmLargeItem(); @@ -1508,7 +1543,8 @@ typename NvmCache::WriteHandle NvmCache::createItem( } XDCHECK(chainedIt->isChainedItem()); XDCHECK_LE(cBlob.data.size(), getStorageSizeInNvm(*chainedIt)); - ::memcpy(chainedIt->getMemory(), cBlob.data.data(), cBlob.data.size()); + copyBlobData(chainedIt->getMemory(), cBlob.data.data(), + cBlob.data.size()); cache_.addChainedItem(it, std::move(chainedIt)); XDCHECK(it->hasChainedItem()); } @@ -1523,6 +1559,22 @@ typename NvmCache::WriteHandle NvmCache::createItem( return it; } +template +void NvmCache::copyBlobData(void* dest, const void* src, size_t n) { + if (copyOutOffload_ && n >= config_.copyOutOffloadMinSize) { + if (navy::copyWithOffload(reinterpret_cast(dest), + reinterpret_cast(src), n)) { + stats().numNvmCopyOutOffloaded.inc(); + stats().numNvmCopyOutOffloadedBytes.add(n); + return; + } + // copyWithOffload already completed the copy on the CPU. + stats().numNvmCopyOutFallbacks.inc(); + return; + } + ::memcpy(dest, src, n); +} + template std::unique_ptr NvmCache::createItemAsIOBuf( folly::StringPiece key, const NvmItem& nvmItem, bool parentOnly) { @@ -1558,7 +1610,7 @@ std::unique_ptr NvmCache::createItemAsIOBuf( XDCHECK_LE(pBlob.origAllocSize, pBlob.data.size()); if (!useCustomCb) { - ::memcpy(item->getMemory(), pBlob.data.data(), pBlob.data.size()); + copyBlobData(item->getMemory(), pBlob.data.data(), pBlob.data.size()); } item->markNvmClean(); @@ -1592,8 +1644,8 @@ std::unique_ptr NvmCache::createItemAsIOBuf( // Propagate the payload directly from Blob only if no customized callback // is set. if (!useCustomCb) { - ::memcpy(chainedItem->getMemory(), cBlob.data.data(), - cBlob.origAllocSize); + copyBlobData(chainedItem->getMemory(), cBlob.data.data(), + cBlob.origAllocSize); } head->appendChain(std::move(chained)); item->markHasChainedItem(); diff --git a/cachelib/cachebench/cache/Cache.h b/cachelib/cachebench/cache/Cache.h index 7ccbf35cee..faa12c81d4 100644 --- a/cachelib/cachebench/cache/Cache.h +++ b/cachelib/cachebench/cache/Cache.h @@ -30,6 +30,7 @@ #include #include #include +#include #include #include #include @@ -770,6 +771,14 @@ Cache::Cache(const CacheConfig& config, auto& bcConfig = nvmConfig.navyConfig.blockCache() .setDataChecksum(config_.navyDataChecksum) + .setChecksumOffload(config_.navyChecksumOffload, + config_.navyChecksumOffloadMinSize) + .setChecksumOffloadReadMinSize( + config_.navyChecksumOffloadReadMinSize) + .setChecksumOffloadCacheControl( + config_.navyChecksumOffloadCacheControl) + .setDirectFlush(config_.navyBlockCacheDirectFlush) + .setFlushCopyOffload(config_.navyBlockCacheFlushCopyOffload) .setCleanRegions(config_.navyCleanRegions, config_.navyCleanRegionThreads) .setRegionSize(config_.navyRegionSizeMB * MB) @@ -811,7 +820,8 @@ Cache::Cache(const CacheConfig& config, .setSizePctAndMaxItemSize(config_.navyBigHashSizePct, config_.navySmallItemMaxSize) .setBucketSize(config_.navyBigHashBucketSize) - .setBucketBfSize(config_.navyBloomFilterPerBucketSize); + .setBucketBfSize(config_.navyBloomFilterPerBucketSize) + .setChecksumOffload(config_.navyBigHashChecksumOffload); } const auto numArenas = config_.getNavyNumArenas(); @@ -1668,11 +1678,18 @@ void Cache::setStringItem(WriteHandle& handle, } auto ptr = reinterpret_cast(getMemory(handle)); - std::strncpy(ptr, str.c_str(), dataSize); + // memcpy/memset instead of strncpy: strncpy scans for the terminator byte + // by byte and is not interposed by DTO, while memcpy/memset of large + // values can be offloaded to DSA. Like strncpy, write exactly dataSize + // bytes: the string (truncated if needed), then a zero-filled tail. + const size_t copyLen = std::min(str.size() + 1, dataSize); + std::memcpy(ptr, str.c_str(), copyLen); // Make sure the copied string ends with null char if (str.size() + 1 > dataSize) { ptr[dataSize - 1] = '\0'; + } else if (copyLen < dataSize) { + std::memset(ptr + copyLen, 0, dataSize - copyLen); } } diff --git a/cachelib/cachebench/cache/components/FlashComponent.cpp b/cachelib/cachebench/cache/components/FlashComponent.cpp index f0a06b80c7..cf52cd0e01 100644 --- a/cachelib/cachebench/cache/components/FlashComponent.cpp +++ b/cachelib/cachebench/cache/components/FlashComponent.cpp @@ -144,6 +144,13 @@ std::unique_ptr createFlashCacheComponent( bcConfig.indexConfig.enableTrackItemHistory(); } bcConfig.indexConfig.validate(); + bcConfig.checksumOffload = config.navyChecksumOffload; + bcConfig.checksumOffloadMinSize = config.navyChecksumOffloadMinSize; + bcConfig.checksumOffloadReadMinSize = config.navyChecksumOffloadReadMinSize; + bcConfig.checksumOffloadCacheControl = + config.navyChecksumOffloadCacheControl; + bcConfig.directFlush = config.navyBlockCacheDirectFlush; + bcConfig.flushCopyOffload = config.navyBlockCacheFlushCopyOffload; // Note: LRU is not yet supported if (config.navySegmentedFifoSegmentRatio.empty() || config.navySegmentedFifoSegmentRatio.size() == 1) { diff --git a/cachelib/cachebench/util/CacheConfig.cpp b/cachelib/cachebench/util/CacheConfig.cpp index 3e87231bcc..b0e5906ef0 100644 --- a/cachelib/cachebench/util/CacheConfig.cpp +++ b/cachelib/cachebench/util/CacheConfig.cpp @@ -116,6 +116,13 @@ CacheConfig::CacheConfig(const folly::dynamic& configJson) { JSONSetVal(configJson, navyAdmissionWriteRateMB); JSONSetVal(configJson, navyMaxConcurrentInserts); JSONSetVal(configJson, navyDataChecksum); + JSONSetVal(configJson, navyChecksumOffload); + JSONSetVal(configJson, navyBlockCacheDirectFlush); + JSONSetVal(configJson, navyBlockCacheFlushCopyOffload); + JSONSetVal(configJson, navyChecksumOffloadMinSize); + JSONSetVal(configJson, navyChecksumOffloadReadMinSize); + JSONSetVal(configJson, navyChecksumOffloadCacheControl); + JSONSetVal(configJson, navyBigHashChecksumOffload); JSONSetVal(configJson, truncateItemToOriginalAllocSizeInNvm); JSONSetVal(configJson, navyEncryption); JSONSetVal(configJson, deviceMaxWriteSize); @@ -193,7 +200,7 @@ CacheConfig::CacheConfig(const folly::dynamic& configJson) { // if you added new fields to the configuration, update the JSONSetVal // to make them available for the json configs and increment the size // below - checkCorrectSize(); + checkCorrectSize(); if (numPools != poolSizes.size()) { throw std::invalid_argument(fmt::format( diff --git a/cachelib/cachebench/util/CacheConfig.h b/cachelib/cachebench/util/CacheConfig.h index 8057399030..ffacea050e 100644 --- a/cachelib/cachebench/util/CacheConfig.h +++ b/cachelib/cachebench/util/CacheConfig.h @@ -271,6 +271,26 @@ struct CacheConfig : public JSONConfig { // default bool navyDataChecksum{true}; + // offloads navy BlockCache data checksumming (fused with the value copy on + // the write path) to Intel DSA via the DTO library, for values of at least + // navyChecksumOffloadMinSize bytes. Requires navyDataChecksum and a build + // with BUILD_WITH_DTO. + bool navyChecksumOffload{false}; + // BlockCache flush path: write the region buffer straight to the device + // (no intermediate copy), or do that copy on Intel DSA. + bool navyBlockCacheDirectFlush{false}; + bool navyBlockCacheFlushCopyOffload{false}; + uint32_t navyChecksumOffloadMinSize{16384}; + // separate gate for read-side verification (0 = same as the write gate; + // 4294967295 disables read-side offload), and the cache-control hint on the + // fused write copy. See BlockCache::Config. + uint32_t navyChecksumOffloadReadMinSize{0}; + bool navyChecksumOffloadCacheControl{true}; + + // offloads navy BigHash bucket checksumming to Intel DSA. Most beneficial + // with large buckets (navyBigHashBucketSize of 16KB+). + bool navyBigHashChecksumOffload{false}; + // by default, only store the size requested by the user into nvm cache bool truncateItemToOriginalAllocSizeInNvm = false; diff --git a/cachelib/cmake/cachelib-config.cmake.in b/cachelib/cmake/cachelib-config.cmake.in index 6e850f5561..a6a6475ce1 100644 --- a/cachelib/cmake/cachelib-config.cmake.in +++ b/cachelib/cmake/cachelib-config.cmake.in @@ -35,6 +35,13 @@ find_dependency(fizz ) find_dependency(wangle) find_dependency(FBThrift) +# When cachelib was built with DSA checksum offload, cachelib_navy links the +# DTO runtime; resolve its imported target for consumers. A DTO fetched and +# built in-tree installs its package into this same prefix. +if (@BUILD_WITH_DTO@) + find_dependency(DTO CONFIG) +endif() + if (NOT TARGET cachelib) include("${CACHELIB_CMAKE_DIR}/cachelib-targets.cmake") endif() diff --git a/cachelib/common/CMakeLists.txt b/cachelib/common/CMakeLists.txt index 1f6c5a0e75..c492cda906 100644 --- a/cachelib/common/CMakeLists.txt +++ b/cachelib/common/CMakeLists.txt @@ -51,6 +51,15 @@ target_link_libraries(cachelib_common PUBLIC ${XXHASH_LIBRARY} ) +if (BUILD_WITH_DTO) + # PUBLIC so the DTO headers/link dependency and the compile definition + # propagate to every cachelib library and to installed-package consumers. + # DTO itself is resolved (installed package or pinned fetch) in the + # top-level CMakeLists.txt. + target_link_libraries(cachelib_common PUBLIC DTO::dto) + target_compile_definitions(cachelib_common PUBLIC CACHELIB_BUILD_WITH_DTO) +endif() + install(TARGETS cachelib_common EXPORT cachelib-exports DESTINATION ${LIB_INSTALL_DIR} ) diff --git a/cachelib/navy/CMakeLists.txt b/cachelib/navy/CMakeLists.txt index e4d9af3eec..ec85134f00 100644 --- a/cachelib/navy/CMakeLists.txt +++ b/cachelib/navy/CMakeLists.txt @@ -32,6 +32,7 @@ add_library (cachelib_navy block_cache/RegionManager.cpp block_cache/SparseMapIndex.cpp common/Buffer.cpp + common/ChecksumOffload.cpp common/Device.cpp common/FdpNvme.cpp common/Hash.cpp @@ -65,6 +66,10 @@ target_link_libraries(cachelib_navy PUBLIC GTest::gmock ) +# DTO (DSA offload runtime) is resolved in the top-level CMakeLists.txt and +# linked PUBLIC into cachelib_common, so cachelib_navy inherits DTO::dto and +# the CACHELIB_BUILD_WITH_DTO definition transitively. + install(TARGETS cachelib_navy EXPORT cachelib-exports DESTINATION ${LIB_INSTALL_DIR} ) @@ -96,6 +101,7 @@ if (BUILD_TESTS) endfunction() add_source_test (common/tests/BufferTest.cpp) + add_source_test (common/tests/ChecksumOffloadTest.cpp) add_source_test (common/tests/HashTest.cpp) add_source_test (common/tests/UtilsTest.cpp) add_source_test (bighash/tests/BucketStorageTest.cpp) diff --git a/cachelib/navy/Factory.cpp b/cachelib/navy/Factory.cpp index e6a2632eba..5146f8936a 100644 --- a/cachelib/navy/Factory.cpp +++ b/cachelib/navy/Factory.cpp @@ -18,6 +18,7 @@ #include #include +#include #include @@ -27,6 +28,7 @@ #include "cachelib/navy/block_cache/BlockCache.h" #include "cachelib/navy/block_cache/FifoPolicy.h" #include "cachelib/navy/block_cache/LruPolicy.h" +#include "cachelib/navy/common/ChecksumOffload.h" #include "cachelib/navy/common/Device.h" #include "cachelib/navy/driver/Driver.h" @@ -37,6 +39,31 @@ namespace facebook::cachelib::navy { namespace { +// Returns true if checksum offload may be enabled. Logs and returns false +// (callers fall back to software checksums) when built without DTO support +// or when the runtime DSA-vs-CPU checksum parity self-check fails. The +// self-check protects against a mixed deployment where accelerator-computed +// checksums would not verify against CPU-computed ones. +bool verifyChecksumOffload(folly::StringPiece engine) { + if (!checksumOffloadSupported()) { + XLOGF(WARN, + "{}: checksum offload requested, but this binary was built without " + "DTO support. Falling back to software checksums.", + engine); + return false; + } + if (!checksumOffloadSelfCheck()) { + XLOGF(ERR, + "{}: DSA checksum offload self-check failed (no usable DSA work " + "queue - descriptors did not complete on the device - or the " + "accelerator checksum does not match the CPU checksum). Falling " + "back to software checksums.", + engine); + return false; + } + return true; +} + class BlockCacheProtoImpl final : public BlockCacheProto { public: BlockCacheProtoImpl() = default; @@ -56,6 +83,19 @@ class BlockCacheProtoImpl final : public BlockCacheProto { void setChecksum(bool enable) override { config_.checksum = enable; } + void setChecksumOffload(bool enable, uint32_t minSize) override { + config_.checksumOffload = enable; + config_.checksumOffloadMinSize = minSize; + } + + void setChecksumOffloadReadMinSize(uint32_t minSize) override { + config_.checksumOffloadReadMinSize = minSize; + } + + void setChecksumOffloadCacheControl(bool enable) override { + config_.checksumOffloadCacheControl = enable; + } + void setLruEvictionPolicy() override { if (!(config_.cacheSize > 0 && config_.regionSize > 0)) { throw std::logic_error("layout is not set"); @@ -136,6 +176,9 @@ class BlockCacheProtoImpl final : public BlockCacheProto { } void setDirectFlush(bool enable) override { config_.directFlush = enable; } + void setFlushCopyOffload(bool enable) override { + config_.flushCopyOffload = enable; + } void setUseCombinedEntryBlock(bool useCombinedEntryBlock) override { config_.useCombinedEntryBlock = useCombinedEntryBlock; @@ -155,6 +198,9 @@ class BlockCacheProtoImpl final : public BlockCacheProto { DestructorCallback cb) && { config_.checkExpired = std::move(checkExpired); config_.destructorCb = std::move(cb); + if (config_.checksumOffload && !verifyChecksumOffload("BlockCache")) { + config_.checksumOffload = false; + } return std::make_unique(std::move(config_)); } @@ -189,6 +235,10 @@ class BigHashProtoImpl final : public BigHashProto { config_.numMutexesPower = numMutexesPower; } + void setChecksumOffload(bool enable) override { + config_.checksumOffload = enable; + } + void setDevice(Device* device) { config_.device = device; } void setDestructorCb(DestructorCallback cb) { @@ -197,6 +247,9 @@ class BigHashProtoImpl final : public BigHashProto { std::unique_ptr create(ExpiredCheck checkExpired) && { config_.checkExpired = std::move(checkExpired); + if (config_.checksumOffload && !verifyChecksumOffload("BigHash")) { + config_.checksumOffload = false; + } if (bloomFilterEnabled_) { if (config_.bucketSize == 0) { throw std::invalid_argument{"invalid bucket size"}; diff --git a/cachelib/navy/Factory.h b/cachelib/navy/Factory.h index 0d487c8b8c..4f7e2e8ba7 100644 --- a/cachelib/navy/Factory.h +++ b/cachelib/navy/Factory.h @@ -50,6 +50,17 @@ class BlockCacheProto { // Enable data checksumming (default: disabled) virtual void setChecksum(bool enable) = 0; + // (Optional) Offload data checksumming to Intel DSA (fused with the value + // copy on the write path) for values of at least @minSize bytes. Requires + // checksumming enabled and a build with DTO support. + virtual void setChecksumOffload(bool enable, uint32_t minSize) = 0; + + // (Optional) Separate size gate for read-side checksum verification; + // 0 = same as the write gate. (Optional) Cache-control hint on the fused + // write-path copy (default on). See BlockCache::Config. + virtual void setChecksumOffloadReadMinSize(uint32_t minSize) = 0; + virtual void setChecksumOffloadCacheControl(bool enable) = 0; + // set*EvictionPolicy function family: sets eviction policy. Supports LRU, // LRU with deferred insert and FIFO. Must set up one of them. @@ -97,6 +108,8 @@ class BlockCacheProto { // (Optional) Set if direct flush without intermediate copy is enabled. virtual void setDirectFlush(bool enable) = 0; + // Copy the region buffer into the flush write buffer on Intel DSA + virtual void setFlushCopyOffload(bool enable) = 0; // (Optional) Set if the combined entry block is enabled. virtual void setUseCombinedEntryBlock(bool useCombinedEntryBlock) = 0; @@ -135,6 +148,10 @@ class BigHashProto { // (Optional) Set number of mutexes for bucket locking as power of 2. // numMutexes = 1 << numMutexesPower. Default: 14 (16K mutexes) virtual void setNumMutexesPower(uint8_t numMutexesPower) = 0; + + // (Optional) Offload bucket checksumming to Intel DSA. Requires a build + // with DTO support. + virtual void setChecksumOffload(bool enable) = 0; }; class EnginePairProto { diff --git a/cachelib/navy/bighash/BigHash.cpp b/cachelib/navy/bighash/BigHash.cpp index e7a1efc17f..687af5b37f 100644 --- a/cachelib/navy/bighash/BigHash.cpp +++ b/cachelib/navy/bighash/BigHash.cpp @@ -86,6 +86,7 @@ BigHash::BigHash(Config&& config, ValidConfigTag) } }}, bucketSize_{config.bucketSize}, + checksumOffload_{config.checksumOffload}, cacheBaseOffset_{config.cacheBaseOffset}, numBuckets_{config.numBuckets()}, bloomFilter_{std::move(config.bloomFilter)}, @@ -366,16 +367,17 @@ Status BigHash::insert(HashedKey hk, bucket->insert(hk, value, checkExpired_, cb); newRemainingBytes = bucket->remainingBytes(); - // rebuild / fix the bloom filter before we move the buffer to do the - // actual write - if (removed + evicted == 0) { - // In case nothing was removed or evicted, we can just add - bfSet(bid, hk.keyHash()); - } else { - bfRebuild(bid, bucket); - } - - const auto res = writeBucket(bid, std::move(buffer)); + // The bloom filter update runs as overlap work: concurrently with the + // bucket checksum when offloaded to DSA, right before it otherwise. In + // both cases it completes before the bucket is written to the device. + const auto res = writeBucket(bid, std::move(buffer), [&] { + if (removed + evicted == 0) { + // In case nothing was removed or evicted, we can just add + bfSet(bid, hk.keyHash()); + } else { + bfRebuild(bid, bucket); + } + }); if (!res) { bfClear(bid); ioErrorCount_.inc(); @@ -512,11 +514,11 @@ Status BigHash::remove(HashedKey hk) { } newRemainingBytes = bucket->remainingBytes(); - // We compute bloom filter before writing the bucket because when encryption - // is enabled, we will "move" the bucket content into writeBucket(). - bfRebuild(bid, bucket); - - const auto res = writeBucket(bid, std::move(buffer)); + // The bloom filter rebuild runs as overlap work: concurrently with the + // bucket checksum when offloaded to DSA, right before it otherwise. In + // both cases it completes before the bucket is written to the device. + const auto res = + writeBucket(bid, std::move(buffer), [&] { bfRebuild(bid, bucket); }); if (!res) { bfClear(bid); ioErrorCount_.inc(); @@ -606,9 +608,12 @@ Buffer BigHash::readBucket(BucketId bid) { auto* bucket = reinterpret_cast(buffer.data()); - auto checksumCheck = [](auto* b, auto bufferView) { - const bool checksumSuccess = - Bucket::computeChecksum(bufferView) == b->getChecksum(); + auto checksumCheck = [this](auto* b, auto bufferView) { + const uint32_t computed = checksumOffload_ + ? Bucket::computeChecksumOffloaded( + bufferView, [] {}) + : Bucket::computeChecksum(bufferView); + const bool checksumSuccess = computed == b->getChecksum(); // TODO (T93631284) we only read a bucket if the bloom filter indicates that // the bucket could have the element. Hence, if check sum errors out and // bloom filter is enabled, we could record the checksum error. However, @@ -632,9 +637,17 @@ Buffer BigHash::readBucket(BucketId bid) { return buffer; } -bool BigHash::writeBucket(BucketId bid, Buffer buffer) { +bool BigHash::writeBucket(BucketId bid, + Buffer buffer, + folly::FunctionRef overlap) { auto* bucket = reinterpret_cast(buffer.data()); - bucket->setChecksum(Bucket::computeChecksum(buffer.view())); + if (checksumOffload_) { + bucket->setChecksum( + Bucket::computeChecksumOffloaded(buffer.view(), overlap)); + } else { + overlap(); + bucket->setChecksum(Bucket::computeChecksum(buffer.view())); + } const bool res = device_.write(getBucketOffset(bid), std::move(buffer), placementHandle_); if (!res) { diff --git a/cachelib/navy/bighash/BigHash.h b/cachelib/navy/bighash/BigHash.h index b1b1a2cf60..af8258aabb 100644 --- a/cachelib/navy/bighash/BigHash.h +++ b/cachelib/navy/bighash/BigHash.h @@ -16,6 +16,7 @@ #pragma once +#include #include #include @@ -66,6 +67,12 @@ class BigHash final : public Engine, folly::NonCopyableNonMovable { struct Config { uint32_t bucketSize{4 * 1024}; + // Offload bucket checksumming to Intel DSA via the DTO library when + // available. Most beneficial with large buckets (16KB+), where the bloom + // filter rebuild overlaps with the accelerator computing the checksum. + // No effect when built without DTO support. + bool checksumOffload{false}; + // The range of device that BigHash will access is guaranteed to be // within [baseOffset, baseOffset + cacheSize) uint64_t cacheBaseOffset{}; @@ -187,7 +194,13 @@ class BigHash final : public Engine, folly::NonCopyableNonMovable { BigHash(Config&& config, ValidConfigTag); Buffer readBucket(BucketId bid); - bool writeBucket(BucketId bid, Buffer buffer); + + // Checksums the bucket and writes it to the device. @overlap is CPU work + // (e.g. bloom filter maintenance) run while the checksum is computed by + // DSA when checksum offload is enabled, or immediately before the software + // checksum otherwise. It may read the bucket but must not modify it. + bool writeBucket( + BucketId bid, Buffer buffer, folly::FunctionRef overlap = [] {}); // The corresponding r/w bucket lock must be held during the entire // duration of the read and write operations. For example, during write, @@ -223,11 +236,16 @@ class BigHash final : public Engine, folly::NonCopyableNonMovable { const size_t numMutexes_; // Serialization format version. Never 0. Versions < 10 reserved for testing. - static constexpr uint32_t kFormatVersion = 10; + // Version 11: navy::checksum switched from CRC-32 (IEEE) to CRC-32C + // (Castagnoli) for DSA offload compatibility. Buckets written by prior + // versions would fail checksum verification. + static constexpr uint32_t kFormatVersion = 11; const ExpiredCheck checkExpired_{}; const DestructorCallback destructorCb_{}; const uint64_t bucketSize_{}; + // Whether to offload bucket checksumming to DSA. See Config::checksumOffload. + const bool checksumOffload_{}; const uint64_t cacheBaseOffset_{}; const uint64_t numBuckets_{}; std::unique_ptr bloomFilter_; diff --git a/cachelib/navy/bighash/Bucket.cpp b/cachelib/navy/bighash/Bucket.cpp index 0138396bcb..1927419be6 100644 --- a/cachelib/navy/bighash/Bucket.cpp +++ b/cachelib/navy/bighash/Bucket.cpp @@ -18,6 +18,7 @@ #include +#include "cachelib/navy/common/ChecksumOffload.h" #include "cachelib/navy/common/Hash.h" namespace facebook::cachelib::navy { @@ -53,6 +54,13 @@ uint32_t Bucket::computeChecksum(BufferView view) { return navy::checksum(data); } +uint32_t Bucket::computeChecksumOffloaded(BufferView view, + folly::FunctionRef overlap) { + constexpr auto kChecksumStart = sizeof(checksum_); + auto data = view.slice(kChecksumStart, view.size() - kChecksumStart); + return checksumWithOverlap(data, overlap); +} + Bucket& Bucket::initNew(MutableBufferView view, uint64_t generationTime) { return *new (view.data()) Bucket(generationTime, view.size() - sizeof(Bucket)); diff --git a/cachelib/navy/bighash/Bucket.h b/cachelib/navy/bighash/Bucket.h index eda5ced271..21c6ca565c 100644 --- a/cachelib/navy/bighash/Bucket.h +++ b/cachelib/navy/bighash/Bucket.h @@ -16,6 +16,7 @@ #pragma once +#include #include #include "cachelib/navy/bighash/BucketStorage.h" @@ -70,6 +71,12 @@ class FOLLY_PACK_ATTR Bucket { // User will pass in a view that contains the memory that is a Bucket static uint32_t computeChecksum(BufferView view); + // Same checksum, but computed by DSA (when built with DTO support) with + // @overlap CPU work running while the accelerator operates. @overlap may + // read the bucket but must not modify it. + static uint32_t computeChecksumOffloaded(BufferView view, + folly::FunctionRef overlap); + // Initialize a brand new Bucket given a piece of memory in the case // that the existing bucket is invalid. (I.e. checksum or generation // and generation time for. diff --git a/cachelib/navy/bighash/tests/BigHashTest.cpp b/cachelib/navy/bighash/tests/BigHashTest.cpp index ced952db4d..631fe8c12d 100644 --- a/cachelib/navy/bighash/tests/BigHashTest.cpp +++ b/cachelib/navy/bighash/tests/BigHashTest.cpp @@ -79,6 +79,30 @@ TEST(BigHash, InsertAndRemove) { EXPECT_EQ(Status::NotFound, bh.remove(makeHK("key"))); } +// Exercises the offloaded bucket-checksum path (DSA when built with DTO +// support, software otherwise), including the bloom-filter-rebuild overlap +// in insert/remove and offloaded verification in readBucket. +TEST(BigHash, InsertRemoveChecksumOffload) { + BigHash::Config config; + setLayout(config, 128, 2); + config.checksumOffload = true; + auto device = std::make_unique>(config.cacheSize, 128); + config.device = device.get(); + + BigHash bh(std::move(config)); + + Buffer value; + uint32_t lat = 0; + EXPECT_EQ(Status::Ok, + bh.insert(makeHK("key"), makeView("12345"), 0 /* poolId */, + 0 /* expiryTime */)); + EXPECT_EQ(Status::Ok, bh.lookup(makeHK("key"), value, lat)); + EXPECT_EQ(makeView("12345"), value.view()); + + EXPECT_EQ(Status::Ok, bh.remove(makeHK("key"))); + EXPECT_EQ(Status::NotFound, bh.lookup(makeHK("key"), value, lat)); +} + // without bloom filters, could exist always returns true. TEST(BigHash, CouldExistWithoutBF) { BigHash::Config config; diff --git a/cachelib/navy/block_cache/BlockCache.cpp b/cachelib/navy/block_cache/BlockCache.cpp index a50aee0411..5f0e441352 100644 --- a/cachelib/navy/block_cache/BlockCache.cpp +++ b/cachelib/navy/block_cache/BlockCache.cpp @@ -32,6 +32,7 @@ #include "cachelib/navy/block_cache/HitsReinsertionPolicy.h" #include "cachelib/navy/block_cache/PercentageReinsertionPolicy.h" #include "cachelib/navy/block_cache/SparseMapIndex.h" +#include "cachelib/navy/common/ChecksumOffload.h" #include "cachelib/navy/common/Hash.h" #include "cachelib/navy/common/Types.h" #include "folly/Range.h" @@ -104,6 +105,11 @@ BlockCache::Config& BlockCache::Config::validate() { } } + if (checksumOffload && !checksum) { + throw std::invalid_argument( + "checksum offload requires data checksum to be enabled"); + } + reinsertionConfig.validate(); indexConfig.validate(); @@ -221,6 +227,12 @@ BlockCache::BlockCache(Config&& config, ValidConfigTag) checkExpired_{std::move(config.checkExpired)}, destructorCb_{std::move(config.destructorCb)}, checksumData_{config.checksum}, + checksumOffload_{config.checksumOffload}, + checksumOffloadMinSize_{config.checksumOffloadMinSize}, + checksumOffloadReadMinSize_{config.checksumOffloadReadMinSize + ? config.checksumOffloadReadMinSize + : config.checksumOffloadMinSize}, + checksumOffloadCacheControl_{config.checksumOffloadCacheControl}, device_{*config.device}, allocAlignSize_{calcAllocAlignSize()}, readBufferSize_{config.readBufferSize < kDefReadBufferSize @@ -250,7 +262,8 @@ BlockCache::BlockCache(Config&& config, ValidConfigTag) config.regionManagerFlushAsync, true /* allowReadDuringReclaim */, config.recoverEvictionPolicy, - config.directFlush}, + config.directFlush, + config.flushCopyOffload}, allocator_{regionManager_, config.allocatorsPerPriority}, reinsertionPolicy_{makeReinsertionPolicy(config.reinsertionConfig)} { validate(config); @@ -274,6 +287,15 @@ std::shared_ptr BlockCache::makeReinsertionPolicy( return reinsertionConfig.getCustomPolicy(*index_); } +uint32_t BlockCache::valueChecksum(BufferView value) const { + if (checksumOffload_ && value.size() >= checksumOffloadReadMinSize_) { + // While DSA computes the checksum, the calling fiber yields (or the + // thread pause-polls), freeing the reader for other requests. + return checksumWithOverlap(value, [] {}); + } + return checksum(value); +} + uint32_t BlockCache::serializedSize(uint32_t keySize, uint32_t valueSize) const { uint32_t size = sizeof(EntryDesc) + keySize + valueSize; @@ -557,21 +579,24 @@ std::pair BlockCache::getRandomAlloc(Buffer& value) { } BufferView valueView{desc.valueSize, entryEnd - entrySize}; - if (checksumData_ && desc.cs != checksum(valueView)) { - XLOGF(ERR, - "Item value checksum mismatch in getRandomAlloc(). Region {} is " - "likely corrupted. Expected: {}, Actual: {}, Offset: {}, " - "Physical-offset: {}, Value-size: {}, Payload (hex): {}", - rid.index(), - desc.cs, - checksum(valueView), - addrEnd.offset() - entrySize, - regionManager_.physicalOffset(addrEnd) - entrySize, - desc.valueSize, - // call folly::unhexlify to convert it back to binary data - folly::hexlify( - folly::ByteRange(valueView.data(), valueView.dataEnd()))); - break; + if (checksumData_) { + const uint32_t valueCs = valueChecksum(valueView); + if (desc.cs != valueCs) { + XLOGF(ERR, + "Item value checksum mismatch in getRandomAlloc(). Region {} is " + "likely corrupted. Expected: {}, Actual: {}, Offset: {}, " + "Physical-offset: {}, Value-size: {}, Payload (hex): {}", + rid.index(), + desc.cs, + valueCs, + addrEnd.offset() - entrySize, + regionManager_.physicalOffset(addrEnd) - entrySize, + desc.valueSize, + // call folly::unhexlify to convert it back to binary data + folly::hexlify( + folly::ByteRange(valueView.data(), valueView.dataEnd()))); + break; + } } // confirm that the chosen NvmItem is still being mapped with the key @@ -770,7 +795,7 @@ void BlockCache::onRegionCleanup(RegionId rid, BufferView buffer) { HashedKey hk = makeHK(entryEnd - sizeof(EntryDesc) - desc.keySize, desc.keySize); BufferView value{desc.valueSize, entryEnd - entrySize}; - if (checksumData_ && desc.cs != checksum(value)) { + if (checksumData_ && desc.cs != valueChecksum(value)) { // We do not need to abort here since the EntryDesc checksum was good, so // we can safely proceed to read the next entry. cleanupValueChecksumErrorCount_.inc(); @@ -1022,20 +1047,22 @@ Status BlockCache::onReadCombinedEntryBlock(uint32_t address, } // Check the checksum for the value - if (checksumData_ && entryDesc.cs != checksum(BufferView(entryDesc.valueSize, - buffer.data()))) { - XLOGF(ERR, - "Item value checksum mismatch in onReadCombinedEntryBlock(). " - "Region {} is likely corrupted. Expected: {}, Actual: {}, Offset: " - "{}, Physical-offset: {}, " - "Value-size: {} Payload (hex): {}", - addrEnd.rid().index(), entryDesc.cs, - checksum(BufferView(entryDesc.valueSize, buffer.data())), - addrEnd.offset() - size, - regionManager_.physicalOffset(addrEnd) - size, entryDesc.valueSize, - folly::hexlify(folly::ByteRange( - buffer.data(), buffer.data() + entryDesc.valueSize))); - return Status::ChecksumError; + if (checksumData_) { + const uint32_t cebValueCs = + valueChecksum(BufferView(entryDesc.valueSize, buffer.data())); + if (entryDesc.cs != cebValueCs) { + XLOGF(ERR, + "Item value checksum mismatch in onReadCombinedEntryBlock(). " + "Region {} is likely corrupted. Expected: {}, Actual: {}, Offset: " + "{}, Physical-offset: {}, " + "Value-size: {} Payload (hex): {}", + addrEnd.rid().index(), entryDesc.cs, cebValueCs, + addrEnd.offset() - size, + regionManager_.physicalOffset(addrEnd) - size, entryDesc.valueSize, + folly::hexlify(folly::ByteRange( + buffer.data(), buffer.data() + entryDesc.valueSize))); + return Status::ChecksumError; + } } // shrink and set it as CEB's buffer @@ -1187,23 +1214,26 @@ AllocatorApiResult BlockCache::reinsertOrRemoveItem( } // Validate checksum as we want to reinsert the item - if (checksumData_ && entryDesc.cs != checksum(value)) { - // We do not need to abort here since the EntryDesc checksum was good, so - // we can safely proceed to read the next entry. - XLOGF(ERR, - "Item value checksum mismatch in reinsertOrRemoveItem(). " - "Item is likely corrupted. Item will not be reinserted. " - "We will continue evaluating remaining items in the region." - "Expected: {}, Actual: {}, Value-size: {}, Payload (hex): {}", - entryDesc.cs, - checksum(value), - entryDesc.valueSize, - folly::hexlify(folly::ByteRange(value.data(), value.dataEnd()))); - reclaimValueChecksumErrorCount_.inc(); - removeItem(false); - recordEvent(hk.key(), AllocatorApiEvent::NVM_REINSERT, - AllocatorApiResult::CORRUPTED, entrySize); - return AllocatorApiResult::CORRUPTED; + if (checksumData_) { + const uint32_t reinsertValueCs = valueChecksum(value); + if (entryDesc.cs != reinsertValueCs) { + // We do not need to abort here since the EntryDesc checksum was good, so + // we can safely proceed to read the next entry. + XLOGF(ERR, + "Item value checksum mismatch in reinsertOrRemoveItem(). " + "Item is likely corrupted. Item will not be reinserted. " + "We will continue evaluating remaining items in the region." + "Expected: {}, Actual: {}, Value-size: {}, Payload (hex): {}", + entryDesc.cs, + reinsertValueCs, + entryDesc.valueSize, + folly::hexlify(folly::ByteRange(value.data(), value.dataEnd()))); + reclaimValueChecksumErrorCount_.inc(); + removeItem(false); + recordEvent(hk.key(), AllocatorApiEvent::NVM_REINSERT, + AllocatorApiResult::CORRUPTED, entrySize); + return AllocatorApiResult::CORRUPTED; + } } // Priority of an re-inserted item is determined by its past accesses @@ -1332,20 +1362,44 @@ void BlockCache::writeEntry(MutableBufferView buffer, // Copy descriptor and the key to the end uint8_t* dest = buffer.data() + buffer.size() - sizeof(EntryDesc); - auto desc = new (dest) EntryDesc(static_cast(keySize), - static_cast(value.size()), keyHash, - lastAccessTimeSecs); - - if (checksumData_) { - desc->cs = checksum(value); - } + EntryDesc* desc = nullptr; + + // Constructs the entry descriptor (which computes its own header checksum + // over the fields preceding csSelf) and copies the key in front of it. The + // value area at the head of the buffer is not touched, so this can overlap + // with the value copy. + auto writeDescAndKey = [&] { + desc = new (dest) EntryDesc(static_cast(keySize), + static_cast(value.size()), keyHash, + lastAccessTimeSecs); + if (!combinedEntry) { + // Key info will be inside the value itself for combined entry block + std::memcpy(dest - keySize, hk->key().data(), keySize); + logicalWriteSize = keySize + value.size(); + } + }; - if (!combinedEntry) { - // Key info will be inside the value itself for combined entry block - std::memcpy(dest - keySize, hk->key().data(), keySize); - logicalWriteSize = keySize + value.size(); + if (checksumData_ && checksumOffload_ && + value.size() >= checksumOffloadMinSize_) { + // Fused copy + checksum on DSA. The descriptor and key are written by the + // CPU BEFORE the descriptor is submitted rather than overlapped with it: + // they sit at the tail of the same slot the device writes the value into, + // and when value.size() is not a multiple of a cache line the CPU and the + // device would otherwise store into the same line concurrently (correct + // under coherence, but a false-sharing round trip per insert). The work + // is ~100 ns, so nothing measurable is lost. The value checksum is not + // covered by csSelf, so it can be filled in afterwards. + writeDescAndKey(); + const uint32_t cs = copyAndChecksum(buffer.data(), value, [] {}, + checksumOffloadCacheControl_); + desc->cs = cs; + } else { + writeDescAndKey(); + if (checksumData_) { + desc->cs = checksum(value); + } + std::memcpy(buffer.data(), value.data(), value.size()); } - std::memcpy(buffer.data(), value.data(), value.size()); logicalWrittenCount_.add(logicalWriteSize); } @@ -1429,20 +1483,23 @@ void BlockCache::readEntry(RegionDescriptor& regionDesc, ld.valueSize_ = desc.valueSize; ld.lastAccessTimeSecs_ = desc.lastAccessTimeSecs; auto slice = ld.buffer_.view().slice(0, desc.valueSize); - if (checksumData_ && desc.cs != checksum(slice)) { - XLOG_N_PER_MS(ERR, 10, 10'000) << fmt::format( - "Item value checksum mismatch in readEntry() looking up key {} in " - "Region {}. Expected: {}, Actual: {}, Offset: {}, Physical-offset: {}, " - "Value-size: {} Payload (hex): {}", - key, addrEnd.rid().index(), desc.cs, checksum(slice), - addrEnd.offset() - size, regionManager_.physicalOffset(addrEnd) - size, - slice.size(), - folly::hexlify( - folly::ByteRange(slice.data(), slice.data() + slice.size()))); - ld.buffer_.reset(); - lookupValueChecksumErrorCount_.inc(); - ld.status_ = Status::ChecksumError; - return; + if (checksumData_) { + const uint32_t sliceCs = valueChecksum(slice); + if (desc.cs != sliceCs) { + XLOG_N_PER_MS(ERR, 10, 10'000) << fmt::format( + "Item value checksum mismatch in readEntry() looking up key {} in " + "Region {}. Expected: {}, Actual: {}, Offset: {}, Physical-offset: " + "{}, Value-size: {} Payload (hex): {}", + key, addrEnd.rid().index(), desc.cs, sliceCs, + addrEnd.offset() - size, + regionManager_.physicalOffset(addrEnd) - size, slice.size(), + folly::hexlify( + folly::ByteRange(slice.data(), slice.data() + slice.size()))); + ld.buffer_.reset(); + lookupValueChecksumErrorCount_.inc(); + ld.status_ = Status::ChecksumError; + return; + } } ld.status_ = Status::Ok; } diff --git a/cachelib/navy/block_cache/BlockCache.h b/cachelib/navy/block_cache/BlockCache.h index 41258889d9..6877b25cf0 100644 --- a/cachelib/navy/block_cache/BlockCache.h +++ b/cachelib/navy/block_cache/BlockCache.h @@ -56,6 +56,26 @@ class BlockCache final : public Engine { DestructorCallback destructorCb; // Checksum data read/written bool checksum{}; + // Offload value checksumming (fused with the value copy on the write + // path) to Intel DSA via the DTO library when available. Requires + // checksum to be enabled. No effect when built without DTO support. + bool checksumOffload{false}; + // Minimum value size to use the offloaded (fused) path on the WRITE path; + // smaller values use software copy+checksum. Submitting and polling a + // descriptor costs a few microseconds regardless of size, about what the + // CPU needs to CRC 16-32 KiB; on a size-diverse (CDN) workload a 16 KiB + // gate beat 4 KiB by 4% of process CPU and 32 KiB was indistinguishable. + uint32_t checksumOffloadMinSize{16384}; + // Minimum value size to offload checksum VERIFICATION (lookup, reclaim, + // reinsertion, cleanup). 0 = same as checksumOffloadMinSize. The read side + // has no CPU work to overlap with the accelerator, so its break-even size + // is higher than the fused write; ~UINT32_MAX disables read-side offload. + uint32_t checksumOffloadReadMinSize{0}; + // Cache-control hint on the fused write-path copy: true steers the value + // bytes toward the CPU cache (right when in-memory lookup hits or a CPU + // flush copy read them soon), false keeps them out (right when the next + // reader is the device's DMA engine, e.g. directFlush/flushCopyOffload). + bool checksumOffloadCacheControl{true}; // Base offset and size (in bytes) of cache on the device uint64_t cacheBaseOffset{}; uint64_t cacheSize{}; @@ -116,6 +136,11 @@ class BlockCache final : public Engine { // old behavior for safe rollout. When true, avoids redundant copy. bool directFlush{false}; + // When not directFlush: copy the region buffer into the flush write + // buffer on Intel DSA (DTO batch descriptor) instead of memcpy. Requires + // a DTO build; verified at runtime, falls back to memcpy. + bool flushCopyOffload{false}; + // name of this BC instance std::string name{}; @@ -289,7 +314,10 @@ class BlockCache final : public Engine { private: // Serialization format version. Never 0. Versions < 10 reserved for testing. - static constexpr uint32_t kFormatVersion = 13; + // Version 14: navy::checksum switched from CRC-32 (IEEE) to CRC-32C + // (Castagnoli) for DSA offload compatibility. Entries written by prior + // versions would fail checksum verification. + static constexpr uint32_t kFormatVersion = 14; // This should be at least the nextTwoPow(sizeof(EntryDesc)). static constexpr uint32_t kDefReadBufferSize = 4096; // Default priority for an item inserted into block cache @@ -538,6 +566,18 @@ class BlockCache final : public Engine { const ExpiredCheck checkExpired_; const DestructorCallback destructorCb_; const bool checksumData_{}; + + // Whether to offload value checksum+copy to DSA on the write path, and the + // minimum value size for which to do so. See Config::checksumOffload. + const bool checksumOffload_{}; + const uint32_t checksumOffloadMinSize_{}; + const uint32_t checksumOffloadReadMinSize_{}; + const bool checksumOffloadCacheControl_{}; + + // Computes the checksum of @value for verification (lookup, reclaim and + // cleanup paths), offloading to DSA when checksum offload is enabled and + // the value meets the size gate. + uint32_t valueChecksum(BufferView value) const; // reference to the under-lying device. const Device& device_; // alloc alignment size indicates the granularity of entry sizes on device. diff --git a/cachelib/navy/block_cache/RegionManager.cpp b/cachelib/navy/block_cache/RegionManager.cpp index 5d892c6804..454f1a2908 100644 --- a/cachelib/navy/block_cache/RegionManager.cpp +++ b/cachelib/navy/block_cache/RegionManager.cpp @@ -16,11 +16,42 @@ #include "cachelib/navy/block_cache/RegionManager.h" +#include + +#include + +#include "cachelib/navy/common/ChecksumOffload.h" + #include "cachelib/common/Profiled.h" #include "cachelib/common/inject_pause.h" #include "cachelib/navy/common/Utils.h" namespace facebook::cachelib::navy { +namespace { +// Populate a buffer's pages so that the first write into it - by the CPU or +// by a DMA engine - does not fault. A DSA descriptor that touches an +// unpopulated page stalls on an IOMMU page request, a kernel round trip per +// 4 KiB page that also holds up the other descriptors queued on that engine; +// measured on the BigCache replay as ~100x the per-op accelerator wait and +// +190% CPU when the region buffers came from a fresh mmap. The in-memory +// buffers are all used, so populating them at construction commits nothing +// that would not be committed anyway. +void populateBuffer(Buffer& buf) { + auto* data = buf.data(); + const auto size = buf.size(); +#ifdef MADV_POPULATE_WRITE + if (::madvise(data, size, MADV_POPULATE_WRITE) == 0) { + return; + } +#endif + // Older kernels: touch one byte per page. Buffer contents are undefined + // until written, so the zero is harmless. + for (size_t off = 0; off < size; off += 4096) { + data[off] = 0; + } +} +} // namespace + RegionManager::RegionManager(uint32_t numRegions, uint64_t regionSize, uint64_t baseOffset, @@ -37,7 +68,8 @@ RegionManager::RegionManager(uint32_t numRegions, bool workerFlushAsync, bool allowReadDuringReclaim, bool recoverEvictionPolicy, - bool directFlush) + bool directFlush, + bool flushCopyOffload) : numPriorities_{numPriorities}, inMemBufFlushRetryLimit_{inMemBufFlushRetryLimit}, numRegions_{numRegions}, @@ -51,6 +83,7 @@ RegionManager::RegionManager(uint32_t numRegions, allowReadDuringReclaim_(allowReadDuringReclaim), recoverEvictionPolicy_{recoverEvictionPolicy}, directFlush_{directFlush}, + flushCopyOffload_{flushCopyOffload}, evictCb_{evictCb}, cleanupCb_{cleanupCb}, numInMemBuffers_{numInMemBuffers}, @@ -58,6 +91,20 @@ RegionManager::RegionManager(uint32_t numRegions, XLOGF(INFO, "{} regions, {} bytes each, allowReadDuringReclaim {}, directFlush {}", numRegions_, regionSize_, allowReadDuringReclaim, directFlush); + if (flushCopyOffload_ && directFlush_) { + XLOG(INFO) << "RegionManager: directFlush set, flush copy offload is moot"; + flushCopyOffload_ = false; + } + if (flushCopyOffload_) { + // probe at the size the flush will actually submit + flushCopyOffload_ = copyOffloadSelfCheck(regionSize_); + if (flushCopyOffload_) { + XLOG(INFO) << "RegionManager: flush copy offload to DSA active"; + } else { + XLOG(WARN) << "RegionManager: flush copy offload requested but the " + "DTO/DSA self-check failed; using memcpy"; + } + } for (uint32_t i = 0; i < numRegions; i++) { regions_[i] = std::make_unique(RegionId{i}, regionSize_); } @@ -67,8 +114,17 @@ RegionManager::RegionManager(uint32_t numRegions, for (uint32_t i = 0; i < numInMemBuffers_; i++) { buffers_.push_back( std::make_unique(device.makeIOBuffer(regionSize_))); + populateBuffer(*buffers_.back()); } + // Every flush holds one in-memory buffer, so numInMemBuffers bounds the + // number of write buffers ever needed at once; keeping up to that many + // means a flush burst never allocates (and frees) buffers under a DMA + // engine that may still hold translations for them. + writeBufPoolCap_ = std::max( + {size_t{2}, 2 * static_cast(numWorkers), + static_cast(numInMemBuffers_)}); + for (uint32_t i = 0; i < numWorkers; i++) { auto name = fmt::format("region_manager_{}", i); workers_.emplace_back( @@ -125,6 +181,33 @@ void RegionManager::reset() { resetEvictionPolicy(); } +Buffer RegionManager::acquireWriteBuffer() { + { + std::lock_guard l{writeBufPoolMutex_}; + if (!writeBufPool_.empty()) { + auto buf = std::move(writeBufPool_.back()); + writeBufPool_.pop_back(); + return buf; + } + } + // Pool empty (startup, or more concurrent flushes than the cap): allocate; + // the buffer joins the pool on release if there is room. When the copy into + // it is done by DSA, populate it first (see populateBuffer). + auto buf = device_.makeIOBuffer(regionSize_); + if (flushCopyOffload_) { + populateBuffer(buf); + } + return buf; +} + +void RegionManager::releaseWriteBuffer(Buffer buf) { + std::lock_guard l{writeBufPoolMutex_}; + if (writeBufPool_.size() < writeBufPoolCap_) { + writeBufPool_.push_back(std::move(buf)); + } + // else: drop it; the cap bounds memory kept after a flush burst +} + Region::FlushRes RegionManager::flushBuffer(const RegionId& rid) { auto& region = getRegion(rid); auto callBack = [this](RelAddress addr, BufferView view) { @@ -136,9 +219,34 @@ Region::FlushRes RegionManager::flushBuffer(const RegionId& rid) { return false; } } else { - auto writeBuffer = device_.makeIOBuffer(view.size()); - writeBuffer.copyFrom(0, view); - if (!deviceWrite(addr, std::move(writeBuffer))) { + auto writeBuffer = acquireWriteBuffer(); + XDCHECK_GE(writeBuffer.size(), view.size()); + if (flushCopyOffload_) { + // One 16 MiB Memory Move descriptor (a single descriptor already runs + // at full device speed, ~55 GB/s measured; batching adds nothing and + // DTO's batch path cannot turn cache control off). Cache control off: + // the write buffer is read next by the storage device's DMA. The + // flush fiber sleeps on a timed baton instead of spinning, since this + // thread rarely has other runnable fibers and the copy takes ~1 ms + // under load. Destination pages must be resident (pre-touched buffers) + // or the work queue configured with block-on-fault. + if (copyLargeWithOffload(writeBuffer.data(), view.data(), view.size(), + 1 /* parts */, false /* cacheControl */, + LargeCopyWait::kSleep, 200 /* us */)) { + flushCopyOffloadCount_.inc(); + } else { + flushCopyFallbackCount_.inc(); + } + } else { + writeBuffer.copyFrom(0, view); + } + // Write through the view overload: the buffer stays ours and returns to + // the pool. Device::write(BufferView) copies only when an encryptor is + // configured, so this is not a second copy. + const bool ok = + deviceWrite(addr, BufferView{view.size(), writeBuffer.data()}); + releaseWriteBuffer(std::move(writeBuffer)); + if (!ok) { return false; } } @@ -767,6 +875,39 @@ void RegionManager::getCounters(const CounterVisitor& visitor) const { visitor("navy_bc_external_fragmentation", externalFragmentation_.get()); visitor("navy_bc_physical_written", physicalWrittenCount_.get(), CounterVisitor::CounterType::RATE); + visitor("navy_bc_flush_copy_offloaded", flushCopyOffloadCount_.get(), + CounterVisitor::CounterType::RATE); + visitor("navy_bc_flush_copy_fallbacks", flushCopyFallbackCount_.get(), + CounterVisitor::CounterType::RATE); + { + const auto w = getCopyLargeWaitStats(); + visitor("navy_bc_flush_copy_calls", w.calls); + visitor("navy_bc_flush_copy_yields", w.yields); + visitor("navy_bc_flush_copy_polls", w.polls); + visitor("navy_bc_flush_copy_sleeps", w.sleeps); + const auto c = getChecksumWaitStats(); + visitor("navy_bc_csum_wait_ops", c.calls); + visitor("navy_bc_csum_wait_polls", c.polls); + visitor("navy_bc_csum_wait_us", c.waitUs); + visitor("navy_bc_csum_wait_blocks", c.sleeps); + const auto o = getCopyOutWaitStats(); + visitor("navy_bc_copyout_wait_ops", o.calls); + visitor("navy_bc_copyout_wait_polls", o.polls); + visitor("navy_bc_copyout_wait_us", o.waitUs); + visitor("navy_bc_copyout_wait_blocks", o.sleeps); + visitor("navy_bc_flush_copy_submit_us", w.submitUs); + visitor("navy_bc_flush_copy_wait_us", w.waitUs); + // Device failures redone on the CPU. Non-zero here with the offload + // "active" means descriptors are being rejected (e.g. a work queue + // without block-on-fault) and the accelerator is doing nothing. + visitor("navy_bc_csum_device_fallbacks", c.fallbacks); + visitor("navy_bc_copyout_device_fallbacks", o.fallbacks); + visitor("navy_bc_flush_copy_device_fallbacks", w.fallbacks); + { + std::lock_guard l{writeBufPoolMutex_}; + visitor("navy_bc_flush_writebuf_pooled", writeBufPool_.size()); + } + } visitor("navy_bc_inmem_active", numInMemBufActive_.get()); visitor("navy_bc_inmem_waiting_flush", numInMemBufWaitingFlush_.get()); visitor("navy_bc_inmem_flush_retries", numInMemBufFlushRetries_.get(), diff --git a/cachelib/navy/block_cache/RegionManager.h b/cachelib/navy/block_cache/RegionManager.h index 5dc13befa3..808cc2e12e 100644 --- a/cachelib/navy/block_cache/RegionManager.h +++ b/cachelib/navy/block_cache/RegionManager.h @@ -21,6 +21,7 @@ #include #include +#include #include #include "cachelib/common/AtomicCounter.h" @@ -92,6 +93,9 @@ class RegionManager { // policy ordering across restarts // @param directFlush whether to write region buffer directly // to device without intermediate copy + // @param flushCopyOffload when not directFlush: copy the region + // buffer into the write buffer on Intel DSA + // (DTO batch descriptor) instead of memcpy RegionManager(uint32_t numRegions, uint64_t regionSize, uint64_t baseOffset, @@ -108,7 +112,8 @@ class RegionManager { bool workerFlushAsync, bool allowReadDuringReclaim = false, bool recoverEvictionPolicy = false, - bool directFlush = false); + bool directFlush = false, + bool flushCopyOffload = false); RegionManager(const RegionManager&) = delete; RegionManager& operator=(const RegionManager&) = delete; @@ -383,6 +388,27 @@ class RegionManager { // Whether to write region buffer directly without intermediate copy const bool directFlush_{false}; + // Flush copy on DSA (see ctor); cleared at construction if directFlush is + // set or the DTO/DSA self-check fails. + bool flushCopyOffload_{false}; + mutable AtomicCounter flushCopyOffloadCount_; + mutable AtomicCounter flushCopyFallbackCount_; + + // Pool of device-aligned write buffers for the non-direct flush path. + // Without it every flush allocates a fresh regionSize_ buffer (16 MiB by + // default, ~10k times per 150 GB written): the allocation itself, and for + // the DSA copy fresh pages the device page-faults on. Buffers are handed + // out for the duration of one flush; the pool keeps at most + // max(2 x workers, numInMemBuffers) of them - the most flushes that can be + // in progress at once - so a burst neither allocates under the DMA engine + // nor grows without bound. + Buffer acquireWriteBuffer(); + void releaseWriteBuffer(Buffer buf); + // mutable: getCounters() is const and reports the pool size + mutable std::mutex writeBufPoolMutex_; + std::vector writeBufPool_; + size_t writeBufPoolCap_{0}; + const RegionEvictCallback evictCb_; const RegionCleanupCallback cleanupCb_; diff --git a/cachelib/navy/block_cache/tests/BlockCacheTest.cpp b/cachelib/navy/block_cache/tests/BlockCacheTest.cpp index 416f482cec..1b38e37775 100644 --- a/cachelib/navy/block_cache/tests/BlockCacheTest.cpp +++ b/cachelib/navy/block_cache/tests/BlockCacheTest.cpp @@ -245,6 +245,51 @@ TEST(BlockCache, InsertLookup) { EXPECT_EQ(0, hits[3]); } +// Exercises the fused copy+checksum write path (DSA-offloaded when built +// with DTO support, software otherwise). minSize of 0 forces every insert +// through the fused path; values large and small verify both roundtrip and +// checksum verification on lookup. +TEST(BlockCache, InsertLookupChecksumOffload) { + std::vector log; + + std::vector hits(4); + auto policy = std::make_unique>(&hits); + auto device = createMemoryDevice(kDeviceSize, nullptr /* encryption */); + auto ex = makeJobScheduler(); + auto config = makeConfig(std::move(policy), *device); + config.checksum = true; + config.checksumOffload = true; + config.checksumOffloadMinSize = 0; + auto engine = makeEngine(std::move(config)); + auto driver = makeDriver(std::move(engine), std::move(ex)); + + BufferGen bg; + for (size_t i = 0; i < 16; i++) { + // Mix of sizes: small and near-region-slot-size values + CacheEntry e{bg.gen(8), bg.gen(i % 2 == 0 ? 800 : 3000)}; + EXPECT_EQ(Status::Ok, + driver->insertAsync(e.key(), e.value(), nullptr, 0 /* poolId */, + 0 /* expiryTime */)); + log.push_back(std::move(e)); + } + driver->flush(); + + for (auto& e : log) { + Buffer value; + uint32_t lat = 0; + EXPECT_EQ(Status::Ok, driver->lookup(e.key(), value, lat)); + EXPECT_EQ(e.value(), value.view()); + } + + // No checksum errors must have been recorded + driver->getCounters({[](folly::StringPiece name, double count, + CounterVisitor::CounterType) { + if (name.contains("checksum_error")) { + EXPECT_EQ(0, count) << name; + } + }}); +} + TEST(BlockCache, LookupReturnsLastAccessTime) { auto hits = std::make_unique>(4); auto policy = std::make_unique>(hits.get()); diff --git a/cachelib/navy/common/ChecksumOffload.cpp b/cachelib/navy/common/ChecksumOffload.cpp new file mode 100644 index 0000000000..00007ff2b3 --- /dev/null +++ b/cachelib/navy/common/ChecksumOffload.cpp @@ -0,0 +1,568 @@ +/* + * Copyright (c) Meta Platforms, Inc. and affiliates. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +#include "cachelib/navy/common/ChecksumOffload.h" + +#include +#include +#include +#include + +#include +#include +#include +#include +#include +#include +#include + +#include "cachelib/navy/common/Hash.h" + +#ifdef CACHELIB_BUILD_WITH_DTO +#include +#endif + +namespace facebook { +namespace cachelib { +namespace navy { + +#ifdef CACHELIB_BUILD_WITH_DTO +namespace { +// Storage for the one-shot op used by copyWithOffload(). It lives for the +// duration of one call on the calling thread/fiber; a fiber that yields keeps +// running on the same NavyThread, and only one call is active per fiber at a +// time, but distinct fibers on one thread could interleave, so the storage +// is per call (stack), not thread-local. +struct alignas(64) OpStorage { + unsigned char bytes[sizeof(dto_async_op)]; +}; + +/* True when dto_async_wait()/dto_batch_wait() park the thread rather than + * spin, so handing them the wait actually idles the core. Not cached: + * dto_wait_blocks() reads a global that DTO's lazy init sets on the first + * submit, and a call before that (a stats read, a test) would otherwise pin + * "does not block" for the life of the process. The read is a single load. */ +bool dtoWaitBlocks() { return dto_wait_blocks() != 0; } + +// 6.1/6.2 checksum ops and 6.3 copy-out: how much of the wait is spent spinning +std::atomic gCsumOps{0}, gCsumYields{0}, gCsumPolls{0}, gCsumWaitUs{0}; +std::atomic gCopyOps{0}, gCopyYields{0}, gCopyPolls{0}, gCopyWaitUs{0}; +std::atomic gCsumBlocks{0}, gCopyBlocks{0}, gLargeBlocks{0}; +// Device failures redone on the CPU. Without these the only evidence of a +// dead accelerator is DTO's per-op stderr line: the results stay correct. +std::atomic gCsumFallbacks{0}, gCopyFallbacks{0}, gLargeFallbacks{0}; +} // namespace + +static_assert(sizeof(dto_async_op) <= 192, + "AsyncChecksumOp::opStorage_ too small for dto_async_op"); +static_assert(alignof(dto_async_op) <= 64, + "AsyncChecksumOp::opStorage_ under-aligned for dto_async_op"); + +bool checksumOffloadSupported() { return true; } + +AsyncChecksumOp::~AsyncChecksumOp() { + // A submitted DSA operation writes to opStorage_ (completion record) and, + // for copies, to dest_. It must be drained before this object dies. + if (state_ == State::kDsaPending) { + wait(); + } +} + +void AsyncChecksumOp::submitImpl(Kind kind, + uint8_t* dest, + BufferView src, + bool cacheControl) { + XDCHECK(state_ == State::kIdle); + dest_ = dest; + src_ = src; + kind_ = kind; + + auto* op = reinterpret_cast(opStorage_); + int rc; + switch (kind) { + case Kind::kCopyCrc: + rc = dto_submit_memcpy_crc(op, dest, src.data(), src.size(), cacheControl); + break; + case Kind::kCopy: + rc = dto_submit_memcpy(op, dest, src.data(), src.size(), cacheControl); + break; + case Kind::kCrc: + default: + rc = dto_submit_crc(op, src.data(), src.size()); + break; + } + if (rc == DTO_ASYNC_SUBMITTED) { + state_ = State::kDsaPending; + return; + } + // DTO_ASYNC_FALLBACK: nothing submitted or copied; run on CPU now. + onDevice_ = false; + if (kind != Kind::kCrc) { + std::memcpy(dest, src.data(), src.size()); + } + crc_ = kind == Kind::kCopy ? 0 : checksum(src); + state_ = State::kCpuDone; +} + +uint32_t AsyncChecksumOp::wait() { + XDCHECK(state_ != State::kIdle); + if (state_ == State::kCpuDone) { + state_ = State::kIdle; + return crc_; + } + + auto* op = reinterpret_cast(opStorage_); + // Yielding suspends this fiber and lets other request fibers run on the + // same NavyThread while the accelerator works; that is the actual + // "free the writer thread" mechanism. Yield ONLY when another fiber is + // ready to run: with a lone fiber, yield() returns immediately and the + // loop would busy-spin through the FiberManager at full rate for the + // whole accelerator operation, which is far more expensive than a pause + // poll. Outside fiber context (plain thread-pool schedulers, tests) this + // reduces to a pause-poll. + auto* fm = folly::fibers::onFiber() + ? folly::fibers::FiberManager::getFiberManagerUnsafe() + : nullptr; + int rc; + const auto t0 = std::chrono::steady_clock::now(); + uint64_t yields = 0, polls = 0, blocks = 0; + while ((rc = dto_async_poll(op)) == DTO_ASYNC_PENDING) { + if (fm && fm->hasReadyTasks()) { + // Another request on this thread can run: switching to it is always + // better than giving the core back, and it is what keeps a blocking + // wait from stalling a NavyThread's other fibers. + ++yields; + folly::fibers::yield(); + } else if (dtoWaitBlocks()) { + // Nothing else to run here. Hand the wait to DTO, which parks this + // thread on its completion aggregator and is woken when the device + // finishes. Returns only once the status byte is written, so the + // enclosing poll settles the result on the next turn. + ++blocks; + dto_async_wait(op); + } else { + ++polls; + folly::asm_volatile_pause(); + } + } + gCsumBlocks.fetch_add(blocks, std::memory_order_relaxed); + gCsumOps.fetch_add(1, std::memory_order_relaxed); + gCsumYields.fetch_add(yields, std::memory_order_relaxed); + gCsumPolls.fetch_add(polls, std::memory_order_relaxed); + gCsumWaitUs.fetch_add(std::chrono::duration_cast( + std::chrono::steady_clock::now() - t0).count(), + std::memory_order_relaxed); + state_ = State::kIdle; + if (rc == DTO_ASYNC_DONE) { + onDevice_ = true; + return kind_ == Kind::kCopy + ? 0 + : static_cast(dto_async_crc_val(op)); + } + // Accelerator failure: destination contents are unspecified, so redo the + // whole operation on the CPU. dest_ is not yet visible to readers per the + // submit contract, so overwriting is safe. + onDevice_ = false; + gCsumFallbacks.fetch_add(1, std::memory_order_relaxed); + if (kind_ != Kind::kCrc) { + std::memcpy(dest_, src_.data(), src_.size()); + } + return kind_ == Kind::kCopy ? 0 : checksum(src_); +} + +#else // !CACHELIB_BUILD_WITH_DTO + +bool checksumOffloadSupported() { return false; } + +AsyncChecksumOp::~AsyncChecksumOp() = default; + +void AsyncChecksumOp::submitImpl(Kind kind, + uint8_t* dest, + BufferView src, + bool /* cacheControl */) { + XDCHECK(state_ == State::kIdle); + kind_ = kind; + onDevice_ = false; + if (kind != Kind::kCrc) { + std::memcpy(dest, src.data(), src.size()); + } + crc_ = kind == Kind::kCopy ? 0 : checksum(src); + state_ = State::kCpuDone; +} + +uint32_t AsyncChecksumOp::wait() { + XDCHECK(state_ == State::kCpuDone); + state_ = State::kIdle; + return crc_; +} + +#endif // CACHELIB_BUILD_WITH_DTO + +void AsyncChecksumOp::submitCopyAndChecksum(uint8_t* dest, + BufferView src, + bool cacheControl) { + XDCHECK(dest); + submitImpl(Kind::kCopyCrc, dest, src, cacheControl); +} + +void AsyncChecksumOp::submitChecksum(BufferView src) { + submitImpl(Kind::kCrc, nullptr, src, false); +} + +void AsyncChecksumOp::submitCopy(uint8_t* dest, + BufferView src, + bool cacheControl) { + XDCHECK(dest); + submitImpl(Kind::kCopy, dest, src, cacheControl); +} + +bool copyWithOffload(uint8_t* dest, const uint8_t* src, size_t n) { +#ifdef CACHELIB_BUILD_WITH_DTO + OpStorage storage; + auto* op = reinterpret_cast(storage.bytes); + // Cache control on: the DRAM item is about to be handed to the caller and + // read (served) right away. + const int rc = dto_submit_memcpy(op, dest, src, n, 1 /* cacheControl */); + if (rc != DTO_ASYNC_SUBMITTED) { + std::memcpy(dest, src, n); + return false; + } + auto* fm = folly::fibers::onFiber() + ? folly::fibers::FiberManager::getFiberManagerUnsafe() + : nullptr; + int st; + const auto t0 = std::chrono::steady_clock::now(); + uint64_t yields = 0, polls = 0, blocks = 0; + while ((st = dto_async_poll(op)) == DTO_ASYNC_PENDING) { + if (fm && fm->hasReadyTasks()) { + ++yields; + folly::fibers::yield(); + } else if (dtoWaitBlocks()) { + ++blocks; + dto_async_wait(op); + } else { + ++polls; + folly::asm_volatile_pause(); + } + } + gCopyBlocks.fetch_add(blocks, std::memory_order_relaxed); + gCopyOps.fetch_add(1, std::memory_order_relaxed); + gCopyYields.fetch_add(yields, std::memory_order_relaxed); + gCopyPolls.fetch_add(polls, std::memory_order_relaxed); + gCopyWaitUs.fetch_add(std::chrono::duration_cast( + std::chrono::steady_clock::now() - t0).count(), + std::memory_order_relaxed); + if (st == DTO_ASYNC_DONE) { + return true; + } + gCopyFallbacks.fetch_add(1, std::memory_order_relaxed); + std::memcpy(dest, src, n); + return false; +#else + std::memcpy(dest, src, n); + return false; +#endif +} + +namespace { +std::atomic gCopyLargeYields{0}; +std::atomic gCopyLargePolls{0}; +std::atomic gCopyLargeSleeps{0}; +std::atomic gCopyLargeCalls{0}; +std::atomic gCopyLargeSubmitUs{0}; +std::atomic gCopyLargeWaitUs{0}; +} // namespace + +CopyLargeWaitStats getChecksumWaitStats() { + CopyLargeWaitStats st; + st.yields = gCsumYields.load(); + st.polls = gCsumPolls.load(); + st.sleeps = gCsumBlocks.load(); + st.calls = gCsumOps.load(); + st.waitUs = gCsumWaitUs.load(); + st.fallbacks = gCsumFallbacks.load(); + return st; +} + +CopyLargeWaitStats getCopyOutWaitStats() { + CopyLargeWaitStats st; + st.yields = gCopyYields.load(); + st.polls = gCopyPolls.load(); + st.sleeps = gCopyBlocks.load(); + st.calls = gCopyOps.load(); + st.waitUs = gCopyWaitUs.load(); + st.fallbacks = gCopyFallbacks.load(); + return st; +} + +CopyLargeWaitStats getCopyLargeWaitStats() { + CopyLargeWaitStats st; + st.yields = gCopyLargeYields.load(); + st.polls = gCopyLargePolls.load(); + st.sleeps = gCopyLargeSleeps.load() + gLargeBlocks.load(); + st.calls = gCopyLargeCalls.load(); + st.submitUs = gCopyLargeSubmitUs.load(); + st.waitUs = gCopyLargeWaitUs.load(); + st.fallbacks = gLargeFallbacks.load(); + return st; +} + +bool copyLargeWithOffload(uint8_t* dest, + const uint8_t* src, + size_t n, + size_t parts, + bool cacheControl, + LargeCopyWait wait, + uint32_t sleepPollUs) { +#ifdef CACHELIB_BUILD_WITH_DTO + if (n == 0) { + return true; + } + gCopyLargeCalls.fetch_add(1, std::memory_order_relaxed); + const auto tSubmit = std::chrono::steady_clock::now(); + parts = std::max(1, std::min(parts, DTO_BATCH_MAX)); + // Piece boundaries 4 KiB aligned so the device works on whole pages. + size_t piece = (n + parts - 1) / parts; + piece = (piece + 4095) & ~size_t{4095}; + void* dst[DTO_BATCH_MAX]; + void* srcs[DTO_BATCH_MAX]; + size_t sizes[DTO_BATCH_MAX]; + int count = 0; + for (size_t off = 0; off < n && count < static_cast(DTO_BATCH_MAX); + off += piece, ++count) { + dst[count] = dest + off; + srcs[count] = const_cast(src) + off; + sizes[count] = std::min(piece, n - off); + } + int rc; + dto_batch_op* op = nullptr; + alignas(64) dto_async_op single; + if (count == 1) { + rc = dto_submit_memcpy(&single, dest, src, n, cacheControl ? 1 : 0); + } else { + op = dto_batch_op_new(); + rc = op ? dto_submit_batch_copy(op, dst, srcs, sizes, count) + : DTO_ASYNC_FALLBACK; + if (rc != DTO_ASYNC_SUBMITTED) { + // The device may lack the Batch opcode (DSA 1.0), or the WQ refused the + // descriptor. One Memory Move covers the range on any DSA; try that + // before giving the copy to memcpy. + if (op) { + dto_batch_op_free(op); + op = nullptr; + } + rc = dto_submit_memcpy(&single, dest, src, n, cacheControl ? 1 : 0); + } + } + const auto tWait = std::chrono::steady_clock::now(); + gCopyLargeSubmitUs.fetch_add( + std::chrono::duration_cast(tWait - tSubmit) + .count(), + std::memory_order_relaxed); + if (rc != DTO_ASYNC_SUBMITTED) { + if (op) { + dto_batch_op_free(op); + } + std::memcpy(dest, src, n); + return false; + } + auto* fm = folly::fibers::onFiber() + ? folly::fibers::FiberManager::getFiberManagerUnsafe() + : nullptr; + int st; + uint64_t yields = 0, polls = 0, sleeps = 0, blocks = 0; + while ((st = op ? dto_batch_poll(op) : dto_async_poll(&single)) == + DTO_ASYNC_PENDING) { + if (fm && fm->hasReadyTasks()) { + // Regardless of wait policy: another fiber on this thread can run, and + // parking the thread (dto_*_wait or sleep_for below block the whole + // NavyThread, not just this fiber) would starve it. Only park when + // idle. + ++yields; + folly::fibers::yield(); + } else if (wait == LargeCopyWait::kSleep && dtoWaitBlocks()) { + // DTO parks this thread and the poller wakes it when the device is + // done: no sleep granularity to overshoot, which for a 16 MiB flush + // was most of the measured wait. + ++blocks; + if (op) { + dto_batch_wait(op); + } else { + dto_async_wait(&single); + } + } else if (wait == LargeCopyWait::kSleep) { + // No blocking wait configured and nothing runnable. nanosleep, even on + // a fiber: the FiberManager's timed baton rides the EventBase wheel + // timer whose tick is 10 ms, which made every wait 10 ms. Blocking the + // thread for sleepPollUs delays fibers that become runnable meanwhile by + // at most that much per iteration, and the CPU idles. + ++sleeps; + std::this_thread::sleep_for(std::chrono::microseconds(sleepPollUs)); + } else { + ++polls; + folly::asm_volatile_pause(); + } + } + gCopyLargeYields.fetch_add(yields, std::memory_order_relaxed); + gCopyLargePolls.fetch_add(polls, std::memory_order_relaxed); + gCopyLargeSleeps.fetch_add(sleeps, std::memory_order_relaxed); + gLargeBlocks.fetch_add(blocks, std::memory_order_relaxed); + gCopyLargeWaitUs.fetch_add( + std::chrono::duration_cast( + std::chrono::steady_clock::now() - tWait) + .count(), + std::memory_order_relaxed); + if (op) { + dto_batch_op_free(op); + // dto_batch_poll repairs failed members on the CPU before returning DONE. + return true; + } + if (st == DTO_ASYNC_DONE) { + return true; + } + gLargeFallbacks.fetch_add(1, std::memory_order_relaxed); + std::memcpy(dest, src, n); + return false; +#else + (void)parts; + (void)cacheControl; + (void)wait; + (void)sleepPollUs; + std::memcpy(dest, src, n); + return false; +#endif +} + +bool copyOffloadSelfCheck(size_t size) { + if (!checksumOffloadSupported()) { + return false; + } + const size_t kSize = std::max(size, 4096); + std::vector src(kSize); + std::vector dst(kSize, 0); + std::mt19937 gen{54321}; + for (auto& b : src) { + b = static_cast(gen()); + } + const bool offloaded = copyWithOffload(dst.data(), src.data(), kSize); + if (!offloaded || std::memcmp(dst.data(), src.data(), kSize) != 0) { + XLOGF(WARN, + "copyOffloadSelfCheck: a {}-byte DSA Memory Move did not complete on " + "the device or did not verify (a work queue whose max_transfer_size " + "is below this size rejects it with XFER_ERANGE)", + kSize); + return false; + } + // The flush call site uses a single descriptor (parts = 1); require exactly + // that. A multi-part copy needs the DSA Batch opcode, which DSA 1.0 lacks - + // probe it for the log, but its absence must not disable the feature. + std::fill(dst.begin(), dst.end(), 0); + const bool single = copyLargeWithOffload(dst.data(), src.data(), kSize, 1); + if (!single || std::memcmp(dst.data(), src.data(), kSize) != 0) { + XLOG(WARN) << "copyOffloadSelfCheck: copyLargeWithOffload(parts=1) did not " + "complete on the device or did not verify"; + return false; + } + std::fill(dst.begin(), dst.end(), 0); + const bool batched = copyLargeWithOffload(dst.data(), src.data(), kSize, 4); + const bool batchOk = + batched && std::memcmp(dst.data(), src.data(), kSize) == 0; + // A refused Batch degrades to one descriptor inside copyLargeWithOffload, + // so "OK" here means the multi-part path works, not that the device has + // the Batch opcode. + XLOGF(INFO, + "copyOffloadSelfCheck: {}-byte single-descriptor copy OK on DSA; " + "multi-part path {}", + kSize, batchOk ? "OK (Batch opcode or single-descriptor degrade)" + : "unavailable (single-descriptor copies will be used)"); + return true; +} + +uint32_t copyAndChecksum(uint8_t* dest, + BufferView src, + folly::FunctionRef overlap, + bool cacheControl) { + AsyncChecksumOp op; + op.submitCopyAndChecksum(dest, src, cacheControl); + overlap(); + return op.wait(); +} + +uint32_t checksumWithOverlap(BufferView src, + folly::FunctionRef overlap) { + AsyncChecksumOp op; + op.submitChecksum(src); + overlap(); + return op.wait(); +} + +bool checksumOffloadSelfCheck() { + if (!checksumOffloadSupported()) { + return false; + } + // Exercise both operations on a buffer large enough to exceed DTO's + // minimum-size gates (DTO_CRC_MIN_BYTES / DTO_MIN_BYTES) so the DSA path + // actually runs, and verify parity with navy::checksum() plus copy + // fidelity. If DSA is unavailable, the CPU fallback must also match. + constexpr size_t kSize = 1024 * 1024; + std::vector src(kSize); + std::vector dst(kSize, 0); + std::mt19937 gen{12345}; + for (auto& b : src) { + b = static_cast(gen()); + } + + const BufferView view{src.size(), src.data()}; + const uint32_t sw = checksum(view); + + // A device failure is redone on the CPU and still returns the right + // checksum, so parity alone cannot tell a working accelerator from one that + // rejects every descriptor (e.g. a WQ without block-on-fault). Require that + // both operations actually completed on the device. + AsyncChecksumOp crcOp; + crcOp.submitChecksum(view); + const uint32_t viaCrc = crcOp.wait(); + if (!crcOp.completedOnDevice()) { + XLOG(WARN) << "checksumOffloadSelfCheck: CRC descriptor did not complete " + "on the device (submission refused or device failure); " + "offload would run on the CPU"; + return false; + } + AsyncChecksumOp copyOp; + copyOp.submitCopyAndChecksum(dst.data(), view, true /* cacheControl */); + const uint32_t viaCopy = copyOp.wait(); + if (!copyOp.completedOnDevice()) { + XLOG(WARN) << "checksumOffloadSelfCheck: fused copy+CRC descriptor did not " + "complete on the device; offload would run on the CPU"; + return false; + } + if (viaCrc != sw || viaCopy != sw) { + XLOGF(WARN, + "checksumOffloadSelfCheck: accelerator checksum disagrees with CPU " + "(cpu {:#x} crc {:#x} copy+crc {:#x})", + sw, viaCrc, viaCopy); + return false; + } + if (std::memcmp(dst.data(), src.data(), kSize) != 0) { + XLOG(WARN) << "checksumOffloadSelfCheck: fused copy is not faithful"; + return false; + } + return true; +} + +} // namespace navy +} // namespace cachelib +} // namespace facebook diff --git a/cachelib/navy/common/ChecksumOffload.h b/cachelib/navy/common/ChecksumOffload.h new file mode 100644 index 0000000000..9d96421ea5 --- /dev/null +++ b/cachelib/navy/common/ChecksumOffload.h @@ -0,0 +1,205 @@ +/* + * Copyright (c) Meta Platforms, Inc. and affiliates. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +#pragma once + +#include + +#include "cachelib/navy/common/Buffer.h" + +namespace facebook { +namespace cachelib { +namespace navy { + +// Checksum (and fused copy+checksum) helpers that can offload to Intel DSA +// via the DTO library when built with CACHELIB_BUILD_WITH_DTO. Without DTO +// support (or when DSA is unavailable at runtime), they run equivalent +// software implementations. In all cases the returned checksum is +// navy::checksum() (raw CRC-32C, seed 0) of @src, so offloaded and +// CPU-computed checksums verify against each other interchangeably. + +// A single asynchronous checksum or fused copy+checksum operation. +// submit*() enqueues the operation on DSA and returns immediately (or runs +// it synchronously on the CPU when offload is unavailable); wait() completes +// it. While the accelerator works, wait() yields the calling fiber (when +// running on one, e.g. on a NavyThread) so other requests can execute on +// this thread — this is what makes the offload truly asynchronous. Off +// fibers it polls with a CPU pause. +// +// The object is not copyable or movable: the device holds a pointer into +// its storage while the operation is in flight. Each submit must be paired +// with exactly one wait() before reuse or destruction. +class AsyncChecksumOp { + public: + AsyncChecksumOp() = default; + AsyncChecksumOp(const AsyncChecksumOp&) = delete; + AsyncChecksumOp& operator=(const AsyncChecksumOp&) = delete; + ~AsyncChecksumOp(); + + // Copies @src to @dest and computes checksum(src) as one fused DSA + // operation. @dest must not overlap @src and must not be accessed until + // wait() returns: on a (rare) accelerator failure the copy is redone on + // the CPU, so intermediate destination contents are unspecified. + // @cacheControl directs the DSA copy output toward the CPU cache; use it + // when the destination will be read again soon (e.g. an in-memory region + // buffer that serves lookups and is flushed to the device shortly after). + void submitCopyAndChecksum(uint8_t* dest, BufferView src, bool cacheControl); + + // Computes checksum(src). @src may be read concurrently (e.g. to rebuild + // a bloom filter) but must not be modified until wait() returns. + void submitChecksum(BufferView src); + + // Plain copy of @src to @dest as a DSA Memory Move (no checksum); wait() + // returns 0. Same destination contract as submitCopyAndChecksum. Used for + // the flash-hit copy-out (Navy read buffer -> DRAM item) where Navy has + // already verified the value checksum. + void submitCopy(uint8_t* dest, BufferView src, bool cacheControl); + + // Completes the submitted operation and returns checksum(src). + uint32_t wait(); + + // True iff the most recently waited operation was executed by the + // accelerator. False when it ran on the CPU: submission refused (DTO size + // gate, enqueue failure, non-DTO build) or a device failure redone in + // software. Lets a self-check tell a working accelerator from a silently + // failing one - both return the correct checksum. + bool completedOnDevice() const { return onDevice_; } + + private: + enum class State : uint8_t { kIdle, kCpuDone, kDsaPending }; + // What the submitted operation computes: a checksum only, a fused copy + // plus checksum, or a plain copy. + enum class Kind : uint8_t { kCrc, kCopyCrc, kCopy }; + + void submitImpl(Kind kind, uint8_t* dest, BufferView src, bool cacheControl); + + State state_{State::kIdle}; + Kind kind_{Kind::kCrc}; + bool onDevice_{false}; + uint32_t crc_{0}; + uint8_t* dest_{nullptr}; + BufferView src_; + // Opaque storage for the DTO async operation (descriptor + completion + // record); sized/aligned to hold dto_async_op without exposing dto.h here. + alignas(64) unsigned char opStorage_[192]; +}; + +// Convenience wrappers over AsyncChecksumOp. The @overlap callback is +// invoked exactly once after submission. Use it for CPU work that can +// proceed while the accelerator operates. It may read the source range but +// must not write it, and must not access the destination range at all. +// Note: when the operation falls back to software (non-DTO build, or DSA +// submission failure), the copy+checksum completes synchronously inside the +// submit step, so @overlap runs after the operation rather than overlapping +// it — the callback's constraints above still apply either way. + +// Returns true if this binary was built with DTO/DSA support. +bool checksumOffloadSupported(); + +// Verifies at runtime that the DTO/DSA checksum matches navy::checksum() and +// that the fused copy is faithful. Returns true iff offload is usable and +// consistent; returns false when built without DTO support. Callers should +// enable offload only if this returns true. +bool checksumOffloadSelfCheck(); + +// Copies @src to @dest and returns the checksum of @src, as a single fused +// DSA "Memory Copy with CRC Generation" operation when available. The copy +// lands toward the CPU cache (cache control) since such destinations are +// typically read again soon. +// @cacheControl steers the copy output toward the CPU cache; keep it on when +// the destination is read again by the CPU soon (in-memory lookup hits, a CPU +// flush copy), off when the next reader is a DMA engine. +uint32_t copyAndChecksum(uint8_t* dest, + BufferView src, + folly::FunctionRef overlap, + bool cacheControl = true); + +// Returns the checksum of @src, computed by DSA when available. +uint32_t checksumWithOverlap(BufferView src, + folly::FunctionRef overlap); + +// Plain copy offload (no checksum). copyWithOffload() copies @n bytes from +// @src to @dest, on DSA when the build has DTO support, the size passes DTO's +// gate and submission succeeds, otherwise with memcpy; it returns true iff +// the accelerator performed the copy. The destination pages should be +// resident (or the WQ configured with block-on-fault): a device page fault +// fails the descriptor and the copy is redone on the CPU. Yields the calling +// fiber while the device works, like the checksum ops. +bool copyWithOffload(uint8_t* dest, const uint8_t* src, size_t n); + +// Runtime check that a DSA copy is faithful; false when built without DTO or +// when DSA is unusable. Callers should enable copy offload only if true. +// @size is the largest copy the caller will submit: a work queue's maximum +// transfer size (a WQ attribute, 2 MiB by default on some configurations) +// rejects bigger descriptors with XFER_ERANGE, so a self-check at a smaller +// size would pass and every real copy would then fail on the device. +bool copyOffloadSelfCheck(size_t size = 1024 * 1024); + +// Large copy offload (e.g. a 16 MiB region buffer). With @parts == 1 the +// range is one DSA Memory Move descriptor; with @parts > 1 it is split into +// consecutive pieces submitted as ONE DSA Batch descriptor, so the pieces run +// in parallel on the engines of the work queue's device. Not every device +// implements the Batch opcode (DSA 1.0 does not): if the batch submit is +// refused, the copy degrades to a single descriptor before falling back to +// memcpy. Waits like the checksum ops (yields the fiber when others are +// runnable, else pause-polls). @cacheControl=false keeps the destination out +// of the CPU caches, right for buffers the device DMAs next. +// +// Return value: true iff the copy was submitted to DSA and completed. For +// @parts == 1 that means the device did the whole copy (a device failure is +// redone with memcpy and returns false). For @parts > 1 DTO repairs any piece +// the device failed on the CPU inside its poll and reports only success, so +// "true" attributes the batch to DSA even if some pieces were repaired; use +// parts == 1 where exact attribution matters. Destination pages should be +// resident (pooled, pre-touched buffers) or the WQ block-on-fault. +// How to wait for a large copy. kSpinOrYield is the checksum ops' policy +// (yield to runnable fibers, else pause-poll) and burns the thread for the +// whole device latency when no fiber is runnable. kSleep nanosleeps the +// calling thread for @sleepPollUs between completion checks (the CPU idles; +// the thread's other fibers wait at most one interval); right for flush-type +// copies whose latency is not critical and whose thread has little else +// queued. +enum class LargeCopyWait : uint8_t { kSpinOrYield, kSleep }; + +bool copyLargeWithOffload(uint8_t* dest, + const uint8_t* src, + size_t n, + size_t parts = 1, + bool cacheControl = false, + LargeCopyWait wait = LargeCopyWait::kSpinOrYield, + uint32_t sleepPollUs = 200); + +// Diagnostics for copyLargeWithOffload waits: cumulative fiber yields and +// pause-polls across all calls (process-wide). +struct CopyLargeWaitStats { + uint64_t yields{0}; + uint64_t polls{0}; + uint64_t sleeps{0}; // waits that parked the thread (DTO blocking wait or nanosleep) + uint64_t calls{0}; + uint64_t submitUs{0}; // time in submit (descriptor build + ENQCMD) + uint64_t waitUs{0}; // time waiting for completion + uint64_t fallbacks{0}; // device failures redone on the CPU (silent otherwise) +}; +CopyLargeWaitStats getCopyLargeWaitStats(); +// Same shape for the checksum ops (6.1/6.2) and the flash-hit copy-out (6.3): +// calls, yields, polls and total wait time, to size what a blocking wait +// (e.g. a shared completion poller) could recover. +CopyLargeWaitStats getChecksumWaitStats(); +CopyLargeWaitStats getCopyOutWaitStats(); + +} // namespace navy +} // namespace cachelib +} // namespace facebook diff --git a/cachelib/navy/common/Device.cpp b/cachelib/navy/common/Device.cpp index ee84808fb9..b79f68aa11 100644 --- a/cachelib/navy/common/Device.cpp +++ b/cachelib/navy/common/Device.cpp @@ -1116,6 +1116,15 @@ IoContext* FileDevice::getIoContext() { asyncBase = std::make_unique(qDepthPerContext_, pollMode); } + if (!asyncBase) { + // Reachable when io_uring is requested but this binary was built + // without liburing (CACHELIB_IOURING_DISABLE). Fail loudly instead of + // crashing on a null deref in AsyncIoContext. + throw std::runtime_error( + "async io requested but no AsyncBase backend is available in this " + "build (io_uring support compiled out?)"); + } + auto idx = incrementalIdx_++; tlContext_.reset(new AsyncIoContext(std::move(asyncBase), idx, evb, qDepthPerContext_, useIoUring, diff --git a/cachelib/navy/common/Hash.cpp b/cachelib/navy/common/Hash.cpp index 50ef925e46..67003c2a33 100644 --- a/cachelib/navy/common/Hash.cpp +++ b/cachelib/navy/common/Hash.cpp @@ -18,12 +18,61 @@ #include +#if defined(__x86_64__) +#include +#endif + namespace facebook::cachelib::navy { uint64_t hashBuffer(BufferView key, uint64_t seed) { return folly::hash::SpookyHashV2::Hash64(key.data(), key.size(), seed); } +namespace { +#if defined(__x86_64__) +// Hardware CRC-32C via the SSE4.2 crc32 instruction, compiled with a target +// attribute and dispatched at runtime, so it is used even when folly (whose +// crc32c_hw is gated on compile-time SSE4.2 flags) was built without SIMD +// flags and would silently fall back to a ~20x slower table implementation. +__attribute__((target("sse4.2"))) uint32_t crc32cHardware(const uint8_t* data, + size_t n, + uint32_t crc) { + while (n >= 8) { + uint64_t v; + memcpy(&v, data, 8); + crc = static_cast(_mm_crc32_u64(crc, v)); + data += 8; + n -= 8; + } + if (n >= 4) { + uint32_t v; + memcpy(&v, data, 4); + crc = _mm_crc32_u32(crc, v); + data += 4; + n -= 4; + } + while (n > 0) { + crc = _mm_crc32_u8(crc, *data++); + n--; + } + return crc; +} + +bool hasSse42() { return __builtin_cpu_supports("sse4.2"); } +#endif +} // namespace + uint32_t checksum(BufferView data, uint32_t startingChecksum) { - return folly::crc32(data.data(), data.size(), startingChecksum); + // CRC-32C (Castagnoli) with raw seed semantics (default 0, no final + // inversion). This matches the value produced by chained SSE4.2 + // _mm_crc32_* instructions and by Intel DSA's CRC generation as used by + // the DTO library, so checksums can be offloaded to DSA and verified on + // CPU (and vice versa) interchangeably. +#if defined(__x86_64__) + static const bool kUseHw = hasSse42(); + if (kUseHw) { + return crc32cHardware(data.data(), data.size(), startingChecksum); + } +#endif + return folly::crc32c(data.data(), data.size(), startingChecksum); } } // namespace facebook::cachelib::navy diff --git a/cachelib/navy/common/Hash.h b/cachelib/navy/common/Hash.h index 7eed7464f8..d1c0486208 100644 --- a/cachelib/navy/common/Hash.h +++ b/cachelib/navy/common/Hash.h @@ -27,7 +27,8 @@ namespace navy { // Default hash function uint64_t hashBuffer(BufferView key, uint64_t seed = 0); -// Default checksumming function +// Default checksumming function: CRC-32C (Castagnoli) with raw seed +// semantics (no final inversion), DSA-offload compatible. uint32_t checksum(BufferView data, uint32_t startingChecksum = 0); // Convenience utils to convert a piece of buffer to a hashed key diff --git a/cachelib/navy/common/tests/ChecksumOffloadTest.cpp b/cachelib/navy/common/tests/ChecksumOffloadTest.cpp new file mode 100644 index 0000000000..2514f1519d --- /dev/null +++ b/cachelib/navy/common/tests/ChecksumOffloadTest.cpp @@ -0,0 +1,235 @@ +/* + * Copyright (c) Meta Platforms, Inc. and affiliates. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +#include +#include +#include + +#include +#include +#include + +#include "cachelib/navy/common/ChecksumOffload.h" +#include "cachelib/navy/common/Hash.h" + +namespace facebook::cachelib::navy::tests { +namespace { +std::vector randomBytes(size_t size, uint32_t seed) { + std::mt19937 gen{seed}; + std::vector bytes(size); + for (auto& b : bytes) { + b = static_cast(gen()); + } + return bytes; +} +} // namespace + +// Sizes spanning below and above typical offload gates (DTO_MIN_BYTES and +// navy-level thresholds), so both the DSA path and internal CPU fallbacks +// are exercised when built with DTO support. +const size_t kSizes[] = {1, 13, 175, 4096, 16 * 1024, 64 * 1024, 1024 * 1024}; + +TEST(ChecksumOffload, CopyAndChecksumMatchesSoftware) { + uint32_t seed = 1; + for (auto size : kSizes) { + auto src = randomBytes(size, seed++); + std::vector dst(size, 0xAA); + const BufferView view{src.size(), src.data()}; + + const uint32_t cs = copyAndChecksum(dst.data(), view, [] {}); + EXPECT_EQ(checksum(view), cs) << "size=" << size; + EXPECT_EQ(0, std::memcmp(src.data(), dst.data(), size)) << "size=" << size; + } +} + +TEST(ChecksumOffload, ChecksumWithOverlapMatchesSoftware) { + uint32_t seed = 100; + for (auto size : kSizes) { + auto src = randomBytes(size, seed++); + const BufferView view{src.size(), src.data()}; + EXPECT_EQ(checksum(view), checksumWithOverlap(view, [] {})) + << "size=" << size; + } +} + +TEST(ChecksumOffload, OverlapRunsExactlyOnce) { + for (auto size : {175UL, 1024UL * 1024UL}) { + auto src = randomBytes(size, 7); + std::vector dst(size, 0); + const BufferView view{src.size(), src.data()}; + + int copyOverlapRuns = 0; + copyAndChecksum(dst.data(), view, [&] { copyOverlapRuns++; }); + EXPECT_EQ(1, copyOverlapRuns) << "size=" << size; + + int crcOverlapRuns = 0; + checksumWithOverlap(view, [&] { crcOverlapRuns++; }); + EXPECT_EQ(1, crcOverlapRuns) << "size=" << size; + } +} + +// The overlap callback is allowed to read the source range concurrently with +// the operation (e.g. BigHash rebuilds its bloom filter from the bucket while +// DSA checksums it). +TEST(ChecksumOffload, OverlapMayReadSource) { + const size_t size = 64 * 1024; + auto src = randomBytes(size, 11); + const BufferView view{src.size(), src.data()}; + + uint64_t sum = 0; + const uint32_t cs = checksumWithOverlap(view, [&] { + for (auto b : src) { + sum += b; + } + }); + EXPECT_EQ(checksum(view), cs); + EXPECT_NE(0ULL, sum); +} + +// Multiple async operations in flight from one thread, out-of-order waits. +TEST(ChecksumOffload, AsyncOpsInFlight) { + const size_t size = 64 * 1024; + auto src1 = randomBytes(size, 21); + auto src2 = randomBytes(size, 22); + std::vector dst1(size, 0); + const BufferView v1{src1.size(), src1.data()}; + const BufferView v2{src2.size(), src2.data()}; + + AsyncChecksumOp op1; + AsyncChecksumOp op2; + op1.submitCopyAndChecksum(dst1.data(), v1, true /* cacheControl */); + op2.submitChecksum(v2); + + // Complete in reverse submission order. + EXPECT_EQ(checksum(v2), op2.wait()); + EXPECT_EQ(checksum(v1), op1.wait()); + EXPECT_EQ(0, std::memcmp(src1.data(), dst1.data(), size)); + + // An op object is reusable after wait(). + op1.submitChecksum(v2); + EXPECT_EQ(checksum(v2), op1.wait()); +} + +// wait() on a fiber takes the yield path: other fibers run on the thread +// while the operation is in flight. +TEST(ChecksumOffload, AsyncOpOnFiber) { + const size_t size = 1024 * 1024; + auto src = randomBytes(size, 27); + std::vector dst(size, 0); + const BufferView view{src.size(), src.data()}; + + folly::EventBase evb; + auto& fm = folly::fibers::getFiberManager(evb); + uint32_t got = 0; + bool otherFiberRan = false; + fm.addTask([&] { + AsyncChecksumOp op; + op.submitCopyAndChecksum(dst.data(), view, true /* cacheControl */); + got = op.wait(); + }); + fm.addTask([&] { otherFiberRan = true; }); + evb.loop(); + + EXPECT_EQ(checksum(view), got); + EXPECT_TRUE(otherFiberRan); + EXPECT_EQ(0, std::memcmp(src.data(), dst.data(), size)); +} + +// Destructor drains an in-flight op (does not leave the device writing to +// freed stack memory). +TEST(ChecksumOffload, AsyncOpDrainOnDestroy) { + const size_t size = 256 * 1024; + auto src = randomBytes(size, 33); + std::vector dst(size, 0); + { + AsyncChecksumOp op; + op.submitCopyAndChecksum(dst.data(), BufferView{src.size(), src.data()}, + false /* cacheControl */); + // no wait(): destructor must drain + } + EXPECT_EQ(0, std::memcmp(src.data(), dst.data(), size)); +} + +TEST(ChecksumOffload, SelfCheck) { + if (!checksumOffloadSupported()) { + EXPECT_FALSE(checksumOffloadSelfCheck()); + return; + } + // Built with DTO. The self-check now requires the accelerator to actually + // execute the descriptors (a CPU fallback returns the right checksum and + // used to pass it). Probe first so a host with no usable work queue skips + // rather than fails; a host where the device runs must agree with + // navy::checksum, or offloaded checksums would not verify against CPU ones. + const auto src = randomBytes(1024 * 1024, 12345); + AsyncChecksumOp probe; + probe.submitChecksum(BufferView{src.size(), src.data()}); + probe.wait(); + if (!probe.completedOnDevice()) { + GTEST_SKIP() << "built with DTO but no usable DSA work queue on this host"; + } + EXPECT_TRUE(checksumOffloadSelfCheck()); +} + +TEST(ChecksumOffload, CompletedOnDeviceIsTruthful) { + // The parity of the result cannot distinguish a working accelerator from a + // silently failing one; completedOnDevice() must. + const auto src = randomBytes(1024 * 1024, 777); + std::vector dst(src.size(), 0); + const BufferView view{src.size(), src.data()}; + AsyncChecksumOp op; + op.submitCopyAndChecksum(dst.data(), view, true /* cacheControl */); + EXPECT_EQ(checksum(view), op.wait()); + EXPECT_EQ(0, std::memcmp(dst.data(), src.data(), src.size())); + if (!checksumOffloadSupported()) { + EXPECT_FALSE(op.completedOnDevice()); + } else if (checksumOffloadSelfCheck()) { + // A passing self-check means 1 MiB descriptors run on the device here. + EXPECT_TRUE(op.completedOnDevice()); + } +} + +TEST(ChecksumOffload, CopyLargeDegradesWithoutBatch) { + // parts > 1 needs the DSA Batch opcode, which DSA 1.0 lacks; the copy must + // be correct on every path (batch, single-descriptor degrade, memcpy). + const auto src = randomBytes(1024 * 1024, 779); + for (size_t parts : {size_t{1}, size_t{4}, size_t{64}}) { + std::vector dst(src.size(), 0); + copyLargeWithOffload(dst.data(), src.data(), src.size(), parts); + EXPECT_EQ(0, std::memcmp(dst.data(), src.data(), src.size())) + << "parts=" << parts; + } + if (checksumOffloadSupported() && copyOffloadSelfCheck()) { + // The single-descriptor path the flush uses: reports device completion + // and does not touch the device-fallback counter. + const auto before = getCopyLargeWaitStats().fallbacks; + std::vector dst(src.size(), 0); + EXPECT_TRUE(copyLargeWithOffload(dst.data(), src.data(), src.size(), 1)); + EXPECT_EQ(0, std::memcmp(dst.data(), src.data(), src.size())); + EXPECT_EQ(before, getCopyLargeWaitStats().fallbacks); + } +} + +TEST(ChecksumOffload, WaitStatsExposeFallbacks) { + // Device failures redone on the CPU are counted, and never exceed the + // number of operations that reached the device path. + const auto csum = getChecksumWaitStats(); + const auto copyOut = getCopyOutWaitStats(); + const auto large = getCopyLargeWaitStats(); + EXPECT_LE(csum.fallbacks, csum.calls); + EXPECT_LE(copyOut.fallbacks, copyOut.calls); + EXPECT_LE(large.fallbacks, large.calls); +} +} // namespace facebook::cachelib::navy::tests diff --git a/cachelib/navy/common/tests/HashTest.cpp b/cachelib/navy/common/tests/HashTest.cpp index d69a22bd78..7bc643b25c 100644 --- a/cachelib/navy/common/tests/HashTest.cpp +++ b/cachelib/navy/common/tests/HashTest.cpp @@ -19,6 +19,35 @@ #include "cachelib/navy/common/Hash.h" namespace facebook::cachelib::navy::tests { +// navy::checksum must compute raw CRC-32C (Castagnoli, reflected poly +// 0x82F63B78), seed 0, no final inversion. This is the value produced by +// chained SSE4.2 _mm_crc32_* instructions and by Intel DSA CRC generation +// (as used by the DTO library). If this test fails, DSA-offloaded checksums +// would not verify against CPU-computed ones (and vice versa), and on-disk +// format versions must be revisited. +TEST(Hash, ChecksumIsRawCrc32c) { + auto cs = [](folly::StringPiece s, uint32_t seed = 0) { + return checksum( + BufferView{s.size(), reinterpret_cast(s.data())}, + seed); + }; + EXPECT_EQ(0x00000000u, cs("")); + EXPECT_EQ(0x93AD1061u, cs("a")); + EXPECT_EQ(0x58E3FA20u, cs("123456789")); + EXPECT_EQ(0xB3FEE25Eu, cs("The quick brown fox jumps over the lazy dog")); + + uint8_t bytes[256]; + for (size_t i = 0; i < sizeof(bytes); i++) { + bytes[i] = static_cast(i); + } + EXPECT_EQ(0x2436A9DBu, checksum(BufferView{sizeof(bytes), bytes})); + + // Chaining with a starting checksum must equal a single-shot computation. + folly::StringPiece full{"123456789"}; + auto part1 = cs("1234"); + EXPECT_EQ(cs(full), cs("56789", part1)); +} + TEST(Hash, HashedKeyCollision) { HashedKey hk1{"key 1"}; HashedKey hk2{"key 2"}; diff --git a/run_dsa_cachebench.sh b/run_dsa_cachebench.sh new file mode 100755 index 0000000000..88596eccc1 --- /dev/null +++ b/run_dsa_cachebench.sh @@ -0,0 +1,339 @@ +#!/bin/bash +# Build CacheLib with DTO, configure DSA work queues, and run cachebench +# pinned to socket 0 with DSA checksum offload and (optionally) DTO's +# transparent memcpy/memset interception. +# +# Usage: ./run_dsa_cachebench.sh [options] +# --workload cdn|bigcache +# cdn = synthetic CDN hit-ratio config (default) +# bigcache = production BigCache trace replay +# (~/bigcache_trace_sea1c01_20250414_20250421.json, +# scaled exactly as the offload study: 400GB Navy, +# 16GB parcel, 500k ops/thread, fiber scheduler) +# --config PATH base cachebench config (overrides the workload default) +# --offload on|off|none +# on = Navy data checksums on, computed by DSA +# off = Navy data checksums on, computed in software +# none = base config's own checksum setting +# (cdn: no Navy section added; bigcache: checksums +# off = the study's no-checksum baseline A) +# --intercept on|off transparent DTO interception of large memcpy/memset. +# off (default) sets DTO_MIN_BYTES=1GB so only the explicit +# checksum-offload API uses DSA; on leaves DTO's own +# size threshold (32KB) in effect. +# --fibers on|off Navy fiber scheduler (NavyRequestScheduler, libaio, +# study settings: maxNumReads 1024 / maxNumWrites 1200 / +# qDepth 32). Default off; pass --fibers on to enable. +# --nvm-size-mb N Navy cache size (default: cdn 20480, bigcache 409600) +# --num-ops N override test_config.numOps (per stressor thread; +# bigcache default 500000 = the study's 12M-op replay) +# --page-size 4k|2mb|1gb +# page size backing the DRAM cache (slabs + hash table, +# SysV hugetlb shm). The hugetlb pools must be reserved +# out-of-band (setup_dsa_bench.sh); 4k = default pages. +# --devices "0 2 4 6" DSA device ids to enable (default "0 2 4 6") +# --out DIR output/work directory (default ./dsa_bench_out) +# --skip-build do not (re)build cachebench +# --skip-dsa do not reconfigure DSA work queues +set -euo pipefail + +REPO="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)" +BUILD_DIR="$REPO/build-cachelib" +PREFIX="$REPO/opt/cachelib" +DTO_DIR="${DTO_DIR:-$HOME/DTO}" + +WORKLOAD=cdn +CONFIG="" +OFFLOAD=on +INTERCEPT=off +FIBERS=off +NVM_SIZE_MB="" +NUM_OPS="" +PAGE_SIZE=4k +DEVICES="0 2 4 6" +OUT="$PWD/dsa_bench_out" +SKIP_BUILD=0 +SKIP_DSA=0 + +while [ $# -gt 0 ]; do + case "$1" in + --workload) WORKLOAD="$2"; shift 2 ;; + --config) CONFIG="$2"; shift 2 ;; + --offload) OFFLOAD="$2"; shift 2 ;; + --intercept) INTERCEPT="$2"; shift 2 ;; + --fibers) FIBERS="$2"; shift 2 ;; + --nvm-size-mb) NVM_SIZE_MB="$2"; shift 2 ;; + --num-ops) NUM_OPS="$2"; shift 2 ;; + --page-size) PAGE_SIZE="$2"; shift 2 ;; + --devices) DEVICES="$2"; shift 2 ;; + --out) OUT="$2"; shift 2 ;; + --skip-build) SKIP_BUILD=1; shift ;; + --skip-dsa) SKIP_DSA=1; shift ;; + -h|--help) sed -n '2,33p' "$0"; exit 0 ;; + *) echo "unknown option: $1" >&2; exit 1 ;; + esac +done + +case "$WORKLOAD" in + cdn) + [ -n "$CONFIG" ] || CONFIG="$REPO/cachelib/cachebench/test_configs/hit_ratio/cdn/config.json" + [ -n "$NVM_SIZE_MB" ] || NVM_SIZE_MB=20480 + ;; + bigcache) + [ -n "$CONFIG" ] || CONFIG="$HOME/bigcache_trace_sea1c01_20250414_20250421.json" + [ -n "$NVM_SIZE_MB" ] || NVM_SIZE_MB=409600 + [ -n "$NUM_OPS" ] || NUM_OPS=500000 + ;; + *) echo "--workload must be cdn|bigcache" >&2; exit 1 ;; +esac +case "$OFFLOAD" in on|off|none) ;; *) echo "--offload must be on|off|none" >&2; exit 1 ;; esac +case "$INTERCEPT" in on|off) ;; *) echo "--intercept must be on|off" >&2; exit 1 ;; esac +case "$FIBERS" in on|off) ;; *) echo "--fibers must be on|off" >&2; exit 1 ;; esac +mkdir -p "$OUT" + +export LD_LIBRARY_PATH="${DTO_LIB_DIR:+$DTO_LIB_DIR:}$PREFIX/lib:$PREFIX/lib64:/usr/local/lib:${LD_LIBRARY_PATH:-}" + +# ---------------------------------------------------------------- build ---- +if [ "$SKIP_BUILD" = 0 ]; then + echo "== Building cachebench (BUILD_WITH_DTO=ON) ==" + if [ ! -f "$BUILD_DIR/CMakeCache.txt" ]; then + mkdir -p "$BUILD_DIR" + export CMAKE_PREFIX_PATH="$PREFIX/lib/cmake:$PREFIX:/usr/local${CMAKE_PREFIX_PATH:+:$CMAKE_PREFIX_PATH}" + # CMAKE_DISABLE_FIND_PACKAGE_uring: cachelib enables its io_uring path + # whenever liburing exists on the system, but that requires a folly + # built with liburing; disable the probe so the two cannot disagree. + cmake -S "$REPO/cachelib" -B "$BUILD_DIR" \ + -DCMAKE_INSTALL_PREFIX="$PREFIX" \ + -DCMAKE_BUILD_TYPE=RelWithDebInfo \ + -DBUILD_TESTS=ON \ + -DBUILD_WITH_DTO=ON \ + -DCMAKE_DISABLE_FIND_PACKAGE_uring=ON + fi + cmake --build "$BUILD_DIR" --target cachebench -j "$(nproc)" +fi +CACHEBENCH="$BUILD_DIR/cachebench/cachebench" +[ -x "$CACHEBENCH" ] || { echo "cachebench not found at $CACHEBENCH" >&2; exit 1; } + +# ------------------------------------------------------------ DSA setup ---- +if [ "$SKIP_DSA" = 0 ]; then + echo "== Configuring DSA shared work queues (devices: $DEVICES) ==" + for d in $DEVICES; do + sudo "$DTO_DIR/accelConfig.sh" "$d" yes 0 4 >/dev/null + done + sudo chmod 0666 /dev/dsa/wq*.0 +fi +# only the offload arm needs a device; software arms run anywhere +if [ "$OFFLOAD" = on ]; then + ls /dev/dsa/wq*.0 >/dev/null 2>&1 || { echo "no enabled DSA WQs under /dev/dsa" >&2; exit 1; } + echo "DSA WQs: $(ls /dev/dsa/)" +fi + +# ------------------------------------------------------- derive config ----- +case "$PAGE_SIZE" in + 4k) HUGE_BYTES=0 ;; + 2mb) HUGE_BYTES=2097152 + free=$(cat /sys/kernel/mm/hugepages/hugepages-2048kB/free_hugepages) + [ "$free" -ge 4600 ] || { echo "need >=4600 free 2MB pages (have $free); reserve with setup_dsa_bench.sh" >&2; exit 1; } ;; + 1gb) HUGE_BYTES=1073741824 + free=$(cat /sys/kernel/mm/hugepages/hugepages-1048576kB/free_hugepages) + [ "$free" -ge 26 ] || { echo "need >=26 free 1GB pages (have $free); reserve with setup_dsa_bench.sh" >&2; exit 1; } ;; + *) echo "bad --page-size '$PAGE_SIZE' (4k|2mb|1gb)" >&2; exit 1 ;; +esac +export HUGE_BYTES + +# SysV shm metadata dir; recreated fresh for every run +SHM_CACHE_DIR="$OUT/shm_meta_${WORKLOAD}_page_${PAGE_SIZE}" +export SHM_CACHE_DIR + +RUN_CONFIG="$OUT/config_${WORKLOAD}_offload_${OFFLOAD}_page_${PAGE_SIZE}.json" +CONFIG="$CONFIG" RUN_CONFIG="$RUN_CONFIG" OFFLOAD="$OFFLOAD" \ +NVM_SIZE_MB="$NVM_SIZE_MB" NUM_OPS="$NUM_OPS" OUT="$OUT" FIBERS="$FIBERS" \ +HUGE_BYTES="$HUGE_BYTES" \ +python3 - <<'EOF' +import json, os, shutil, sys +cfgPath = os.environ["CONFIG"] +cfg = json.load(open(cfgPath)) +base = os.path.dirname(os.path.abspath(cfgPath)) +tc = cfg["test_config"] +cc = cfg["cache_config"] +offload = os.environ["OFFLOAD"] +isTraceReplay = "replay" in tc.get("generator", "") + +# cachebench resolves distribution files by prepending the config file's +# directory (even to absolute paths), so copy them next to the derived +# config and reference them by basename. +for k in ("popDistFile", "valSizeDistFile"): + if k in tc: + src = tc[k] if os.path.isabs(tc[k]) else os.path.join(base, tc[k]) + dst = os.path.join(os.environ["OUT"], os.path.basename(src)) + if os.path.abspath(src) != os.path.abspath(dst): + shutil.copyfile(src, dst) + tc[k] = os.path.basename(src) +# the (multi-GB) trace file is opened directly: absolutize, don't copy +if "traceFileName" in tc and not os.path.isabs(tc["traceFileName"]): + tc["traceFileName"] = os.path.join(base, tc["traceFileName"]) +if "traceFileName" in tc and not os.path.exists(tc["traceFileName"]): + sys.exit(f"trace file not found: {tc['traceFileName']}") + +if os.environ["NUM_OPS"]: + tc["numOps"] = int(os.environ["NUM_OPS"]) + +if isTraceReplay: + # Trace replay (BigCache): the base config already has a Navy section + # sized for its original production host; rescale it to this machine + # exactly as the offload study did. + cc["nvmCacheSizeMB"] = int(os.environ["NVM_SIZE_MB"]) + cc["nvmCachePaths"] = [os.path.join(os.environ["OUT"], "navy_cache_file")] + cc["navyParcelMemoryMB"] = 16384 + if offload != "none": + cc["navyDataChecksum"] = True + cc["navyChecksumOffload"] = offload == "on" + # only a default: a base config that sets the gate (a size sweep) wins + cc.setdefault("navyChecksumOffloadMinSize", 16384) + else: + # "no checksum" arm: override whatever the base config says, so a + # config saved from an offload run does not silently keep checksums + # (and a DSA self-check) on + cc["navyDataChecksum"] = False + cc["navyChecksumOffload"] = False + cc["navyBlockCacheFlushCopyOffload"] = False +elif offload != "none": + # Synthetic workloads (CDN): the base config is DRAM-only; add the + # hybrid Navy tier the offload applies to. + cc["nvmCacheSizeMB"] = int(os.environ["NVM_SIZE_MB"]) + cc["nvmCachePaths"] = [os.path.join(os.environ["OUT"], "navy_cache_file")] + # defaults only: a base config may set these (the fiber scheduler needs + # maxNumReads 1024 / maxNumWrites 1200 to divide evenly, so 32 writers + # is invalid with --fibers on; use 30 or 40) + cc.setdefault("navyReaderThreads", 32) + cc.setdefault("navyWriterThreads", 32) + cc["navyDataChecksum"] = True + cc["navyChecksumOffload"] = offload == "on" + cc.setdefault("navyChecksumOffloadMinSize", 16384) + +huge = int(os.environ.get("HUGE_BYTES", "0")) +# DRAM cache (slabs + hash table) on SysV shm so the page-size knob actually +# applies to the cache; hugePageSize 0 = normal 4K pages. Without shmType + +# cacheDir the allocator silently uses heap memory and ignores hugePageSize. +cc["shmType"] = "sysv" +cc["cacheDir"] = os.environ["SHM_CACHE_DIR"] +if huge: + cc["hugePageSize"] = huge + +if os.environ["FIBERS"] == "on": + # NavyRequestScheduler (fibers, libaio) with the study's settings; + # required for the offload's yield-during-DSA-wait to overlap work. + cc["navyMaxNumReads"] = 1024 + cc["navyMaxNumWrites"] = 1200 + cc["navyQDepth"] = 32 + cc["navyEnableIoUring"] = False + +# Print the navy counter map at the end of the run: the offload/fallback/wait +# counters are registered in RegionManager::getCounters but only reach the log +# with this on. A run that cannot be audited from its log is not a measurement. +cc["printNvmCounters"] = True + +json.dump(cfg, open(os.environ["RUN_CONFIG"], "w"), indent=2) +print(f"wrote {os.environ['RUN_CONFIG']}") +EOF + +# ---------------------------------------------------------------- run ------ +# navy with 128 reader + 300 writer threads needs far more fds than the +# usual 1024 soft limit +ulimit -n 65536 2>/dev/null || true + +export DTO_USESTDC_CALLS=0 +export DTO_CRC_MIN_BYTES=4096 +# Keep freed memory mapped. Navy populates its region and flush buffers +# itself (RegionManager), which is what keeps a DSA using shared virtual +# addressing from stalling on IOMMU page requests; these tunables remove the +# remaining allocator churn (a few hundred thousand first-touch page requests +# per run from glibc unmapping and re-obtaining large chunks) and also trim +# the software arm by ~7%, so they are set for every arm. jemalloc users: the +# equivalent is dirty_decay_ms:-1,muzzy_decay_ms:-1. +export GLIBC_TUNABLES="${GLIBC_TUNABLES:-glibc.malloc.mmap_threshold=33554432:glibc.malloc.trim_threshold=4294967295:glibc.malloc.top_pad=67108864}" +export DTO_WAIT_METHOD="${DTO_WAIT_METHOD:-busypoll}" +export DTO_COLLECT_STATS=1 +if [ "$INTERCEPT" = off ]; then + export DTO_MIN_BYTES=1073741824 # 1GB: transparent interception off +else + unset DTO_MIN_BYTES # DTO's own threshold decides +fi + +LOG="$OUT/cachebench_${WORKLOAD}_offload_${OFFLOAD}_intercept_${INTERCEPT}_page_${PAGE_SIZE}.log" +echo "== Running cachebench on socket 0 (workload=$WORKLOAD offload=$OFFLOAD intercept=$INTERCEPT fibers=$FIBERS) ==" +echo " log: $LOG" + +# The fiber scheduler has a known ~1-in-4 startup hang (frozen at +# NavyRequestDispatcher startup). With --progress 60 a healthy run writes to +# its log at least once a minute, so prolonged log silence while the process +# lives means hung: kill and retry (up to 3 attempts). A cold trace preload +# is silent too, so the threshold is configurable (WATCHDOG_SILENCE seconds). +WATCHDOG_SILENCE="${WATCHDOG_SILENCE:-330}" +attempt_run() { + rm -f "$OUT/navy_cache_file" + # PREFILL_NAVY=1: populate the Navy file before the run (tmpfs stores: a + # fresh file is 100 GB of first-touch page allocation inside the measured + # window, which stalls flushes and inflates insert latency for every arm) + if [ -n "${PREFILL_NAVY:-}" ]; then + dd if=/dev/zero of="$OUT/navy_cache_file" bs=1M count="$NVM_SIZE_MB" status=none + fi + rm -rf "$SHM_CACHE_DIR" + # reap orphaned SysV segments (>100MB, ours) from earlier killed runs; + # a clean cachelib exit persists them for warm restart, which would pin + # the hugetlb pool across runs + ipcs -m 2>/dev/null | awk -v u="$USER" '$3==u && $5+0>100000000 {print $2}' | while read -r id; do ipcrm -m "$id" 2>/dev/null || true; done + /usr/bin/time -v numactl -N 0 -m 0 \ + "$CACHEBENCH" --json_test_config "$RUN_CONFIG" --progress 60 --report_api_latency=true \ + > "$LOG" 2>&1 & + local tpid=$! + local last=-1 silent=0 + while kill -0 "$tpid" 2>/dev/null; do + sleep 30 + local sz + sz=$(stat -c %s "$LOG" 2>/dev/null || echo 0) + if [ "$sz" = "$last" ]; then + silent=$((silent + 1)) + if [ "$silent" -ge $((WATCHDOG_SILENCE / 30)) ]; then + pkill -9 -P "$tpid" 2>/dev/null + kill -9 "$tpid" 2>/dev/null + wait "$tpid" 2>/dev/null + return 99 + fi + else + silent=0 + fi + last=$sz + done + wait "$tpid" +} + +set +e +rc=99 +WATCHDOG_RETRIES="${WATCHDOG_RETRIES:-6}" +for attempt in $(seq 1 "$WATCHDOG_RETRIES"); do + attempt_run + rc=$? + [ "$rc" -ne 99 ] && break + echo "WATCHDOG: run frozen for ${WATCHDOG_SILENCE}s+ (startup hang?), retrying (attempt $attempt of $WATCHDOG_RETRIES)" +done +set -e +rm -f "$OUT/navy_cache_file" +rm -rf "$SHM_CACHE_DIR" +ipcs -m 2>/dev/null | awk -v u="$USER" '$3==u && $5+0>100000000 {print $2}' | while read -r id; do ipcrm -m "$id" 2>/dev/null || true; done + +# ------------------------------------------------------------- report ------ +echo "exit=$rc" +grep -E "Total Ops|get |set |NVM Gets|NVM Puts" "$LOG" | head -6 || true +# exclude the printed counter map: its key names (navy_*_checksum_errors : 0) +# would otherwise count as errors +echo "checksum error lines: $(grep -viE '^navy_' "$LOG" | grep -ciE 'checksum.*(error|mismatch)' || true)" +echo "checksum error counters: $(grep -iE '^navy_.*checksum.*error' "$LOG" | sed 's/ : /=/' | tr '\n' ' ')" +echo "-- DSA offload counters (device fallbacks must be 0 for an offload run) --" +grep -E "^navy_bc_(csum|copyout|flush_copy)_|^navy_hugetlb_arena_misses" "$LOG" | sed 's/ : / = /' || true +grep -E "User time|System time|Elapsed \(wall|Maximum resident" "$LOG" || true +if grep -q "Number of Memory Operations" "$LOG"; then + echo "-- DTO op counts (see full table in the log) --" + grep -A20 "Number of Memory Operations" "$LOG" | head -24 +fi +exit $rc