diff --git a/.github/workflows/cuda.yml b/.github/workflows/cuda.yml index dea94bf4ac7..b30da2a99de 100644 --- a/.github/workflows/cuda.yml +++ b/.github/workflows/cuda.yml @@ -515,10 +515,11 @@ jobs: -v -o "addopts=" cmake --preset llm-release-cuda -DEXECUTORCH_BUILD_TESTS=ON - cmake --build cmake-out --target test_cuda_allocator test_cuda_mutable_state test_cuda_weight_cache test_cuda_guard test_cuda_stream_guard -j$(nproc) + cmake --build cmake-out --target test_cuda_allocator test_cuda_mutable_state test_cuda_weight_cache test_cuda_kv_cache test_cuda_guard test_cuda_stream_guard -j$(nproc) ctest --test-dir cmake-out -R test_cuda_allocator --output-on-failure -V ctest --test-dir cmake-out -R test_cuda_mutable_state --output-on-failure -V ctest --test-dir cmake-out -R test_cuda_weight_cache --output-on-failure -V + ctest --test-dir cmake-out -R test_cuda_kv_cache --output-on-failure -V ctest --test-dir cmake-out -R test_cuda_guard --output-on-failure -V ctest --test-dir cmake-out -R test_cuda_stream_guard --output-on-failure -V diff --git a/backends/cuda/CMakeLists.txt b/backends/cuda/CMakeLists.txt index 187fceaab03..b4aae79ba22 100644 --- a/backends/cuda/CMakeLists.txt +++ b/backends/cuda/CMakeLists.txt @@ -351,7 +351,9 @@ set(_aoti_cuda_backend_sources ) # The off-graph KV cache builds on the LLM extension's neutral cache. if(EXECUTORCH_BUILD_EXTENSION_LLM) - list(APPEND _aoti_cuda_backend_sources runtime/cuda_kv_cache.cpp) + list(APPEND _aoti_cuda_backend_sources runtime/cuda_kv_cache.cpp + runtime/cuda_kv_pool.cpp + ) endif() if(_cuda_is_msvc_toolchain) # MSVC links aoti_cuda_backend into portable_lib without relying on C++ @@ -480,4 +482,12 @@ if(BUILD_TESTING) EXTRA_LIBS aoti_cuda_backend ) target_compile_definitions(test_cuda_weight_cache PRIVATE CUDA_AVAILABLE=1) + + if(EXECUTORCH_BUILD_EXTENSION_LLM) + et_cxx_test( + test_cuda_kv_cache SOURCES runtime/test/test_cuda_kv_cache.cpp EXTRA_LIBS + aoti_cuda_backend + ) + target_compile_definitions(test_cuda_kv_cache PRIVATE CUDA_AVAILABLE=1) + endif() endif() diff --git a/backends/cuda/runtime/cuda_kv_cache.cpp b/backends/cuda/runtime/cuda_kv_cache.cpp index 364a27adb8f..a0ade97b598 100644 --- a/backends/cuda/runtime/cuda_kv_cache.cpp +++ b/backends/cuda/runtime/cuda_kv_cache.cpp @@ -8,21 +8,13 @@ #include -#include #include #include #include -#include -#include -#include -#include #include -#include #include -#include -#include -#include +#include #include #include #include @@ -30,56 +22,9 @@ namespace executorch::backends::cuda { namespace { -namespace aoti = ::executorch::backends::aoti; namespace slimc10 = ::executorch::backends::aoti::slim::c10; -using ::executorch::backends::aoti::slim::from_blob; -using ::executorch::backends::aoti::slim::SlimTensor; using ::executorch::runtime::Error; -struct Allocation { - void* k{nullptr}; - void* v{nullptr}; - int64_t rows{0}; -}; - -struct Descriptor { - std::string internal_name; - int64_t layer_id{0}; - bool is_value{false}; -}; - -// The tensors a handle's AOTI constants currently point at. -struct Bound { - std::vector> tensors; -}; - -// Makes the cache's device current for a scope. The backend runs on one -// device, so this is normally a no-op, but event and stream calls act on the -// current device and a caller may have switched it. -class DeviceGuard { - public: - explicit DeviceGuard(int device) { - if (cudaGetDevice(&previous_) == cudaSuccess && previous_ != device) { - restore_ = cudaSetDevice(device) == cudaSuccess; - } - } - ~DeviceGuard() { - if (restore_) { - (void)cudaSetDevice(previous_); - } - } - DeviceGuard(const DeviceGuard&) = delete; - DeviceGuard& operator=(const DeviceGuard&) = delete; - - private: - int previous_{0}; - bool restore_{false}; -}; - -std::string fqn(int64_t layer_id, const char* suffix) { - return "__et_offgraph_kv_layer_" + std::to_string(layer_id) + "_" + suffix; -} - bool is_ring(const cache::LayerGeometry& layer) { return layer.policy.kind == cache::LayerPolicy::Kind::Ring; } @@ -106,13 +51,36 @@ bool storage_dtype_of(int kv_dtype, slimc10::ScalarType& out) { } } +// Slots a ring layer needs to serve one step of up to max_write tokens: the +// step writes all of them before attending, and its earliest query still +// reads back window - 1 positions. Same formula as cache::RingPolicy and as +// ring_physical_capacity() in backends/cuda/passes/lower_offgraph_kv.py. +// make_cuda_sequence_kv_cache() refuses ring layers without max_write. +std::vector sequence_pool_layers( + const cache::CacheGeometry& geometry, + const cache::CacheConfig& cfg) { + std::vector layers; + layers.reserve(geometry.layers.size()); + for (const cache::LayerGeometry& layer : geometry.layers) { + const bool ring = is_ring(layer); + layers.push_back(CudaKVPool::Layer{ + layer.n_kv_heads, + layer.head_dim, + ring ? static_cast(layer.policy.window) + + cfg.max_write.value_or(1) - 1 + : static_cast(cfg.capacity), + /*growable=*/!ring}); + } + return layers; +} + // One sequence's KV storage on one CUDA device. // -// The neutral base owns the logical length, admission and rewind. This class -// owns only bytes: the device allocations, their geometric growth, and the -// AOTI bindings that point the compiled program at them. Physical slot math -// stays in the compiled program, on device, so a captured CUDA graph replays -// against the current positions rather than the ones live at capture. +// The neutral base owns the logical length, admission and rewind; the pool +// owns the bytes and the AOTI bindings that point the compiled program at +// them. Physical slot math stays in the compiled program, on device, so a +// captured CUDA graph replays against the current positions rather than the +// ones live at capture. class CudaSequenceKVCache final : public cache::SequenceCache, public CudaKVCache { public: @@ -122,86 +90,57 @@ class CudaSequenceKVCache final : public cache::SequenceCache, slimc10::ScalarType storage_dtype) : cache::SequenceCache(geometry, cfg), geometry_(geometry), - config_(cfg), - storage_dtype_(storage_dtype) {} + pool_( + sequence_pool_layers(geometry, cfg), + /*side_buffers=*/{}, + storage_dtype, + cfg.initial_capacity) {} CudaSequenceKVCache(const CudaSequenceKVCache&) = delete; CudaSequenceKVCache& operator=(const CudaSequenceKVCache&) = delete; CudaSequenceKVCache(CudaSequenceKVCache&&) = delete; CudaSequenceKVCache& operator=(CudaSequenceKVCache&&) = delete; - // Nothing else can reach the cache once it is being destroyed, so this takes - // no lock. It also never touches stream_: that is the caller's stream, which - // may already be gone. The device is drained instead, which covers any step - // still in flight on whatever stream ran it. - ~CudaSequenceKVCache() override { - if (!device_known_) { - return; - } - DeviceGuard device(device_); - (void)cudaDeviceSynchronize(); - for (auto& allocation : allocations_) { - release(allocation.second, cudaStreamPerThread); - } - (void)cudaStreamSynchronize(cudaStreamPerThread); - if (last_step_done_ != nullptr) { - (void)cudaEventDestroy(last_step_done_); - } - } + ~CudaSequenceKVCache() override = default; // cache::SequenceControl. Keeps the storage: a reset session reuses the // grown allocations, so it neither regrows nor invalidates captured graphs. void clear() override { std::lock_guard guard(mutex_); cache::SequenceCache::clear(); - metrics_.logical_length = 0; ET_LOG( Info, "offgraph_kv: reset flat_capacity=%lld allocated_bytes=%lld", - static_cast(metrics_.flat_capacity), - static_cast(metrics_.allocated_bytes)); + static_cast(pool_.rows()), + static_cast(pool_.allocated_bytes())); } // CudaKVCache. runtime::Result note_handle(CudaDelegateHandle* handle) override { std::lock_guard guard(mutex_); - if (!device_known_) { - if (cudaGetDevice(&device_) != cudaSuccess) { - error_ = Error::Internal; - return error_; - } - device_known_ = true; - } - const Error error = build_descriptors(handle); - if (error != Error::Ok) { - error_ = error; - return error; - } - if (descriptors_[handle].empty()) { - descriptors_.erase(handle); - return false; + auto serves = pool_.note_handle(handle); + if (!serves.ok()) { + error_ = serves.error(); } - handles_associated_ = true; - return true; + return serves; } void forget_handle(CudaDelegateHandle* handle) override { std::lock_guard guard(mutex_); - descriptors_.erase(handle); - bound_.erase(handle); + pool_.forget_handle(handle); } Error rebind_for_execute(CudaDelegateHandle* handle) override { std::lock_guard guard(mutex_); - if (descriptors_.find(handle) == descriptors_.end()) { + if (!pool_.serves(handle)) { return Error::Ok; // a handle with no off-graph storage } ET_CHECK_OK_OR_RETURN_ERROR(error_); ET_CHECK_OR_RETURN_ERROR( - !allocations_.empty(), + pool_.allocated(), InvalidState, "offgraph_kv: prepare_step must run before execute"); - return bind(handle); + return pool_.bind(handle); } Error prepare_step(int64_t write_length, cudaStream_t stream) override { @@ -209,7 +148,10 @@ class CudaSequenceKVCache final : public cache::SequenceCache, ET_CHECK_OK_OR_RETURN_ERROR(check_write_length(write_length)); ET_CHECK_OK_OR_RETURN_ERROR(error_); if (!validated_) { - const Error valid = validate_locked(); + // Whole-model check, so it cannot run until every program has been + // loaded and registered. The first step is the earliest moment that is + // guaranteed. + const Error valid = pool_.validate(); if (valid != Error::Ok) { error_ = valid; return valid; @@ -218,18 +160,7 @@ class CudaSequenceKVCache final : public cache::SequenceCache, } const int position = length(); ET_CHECK_OK_OR_RETURN_ERROR(admit(position, write_length)); - DeviceGuard device(device_); - ET_CHECK_OK_OR_RETURN_ERROR(follow_previous_step(stream)); - stream_ = stream; - const int64_t required = position + write_length; - ET_CHECK_OK_OR_RETURN_ERROR(ensure_initial_allocations(required, stream)); - if (required > metrics_.flat_capacity) { - const int64_t next = std::min( - config_.capacity, - std::max(required, metrics_.flat_capacity * 2)); - ET_CHECK_OK_OR_RETURN_ERROR(grow_flat(next, stream)); - } - return Error::Ok; + return pool_.prepare(position + write_length, position, stream); } Error commit_step(int64_t write_length) override { @@ -244,14 +175,17 @@ class CudaSequenceKVCache final : public cache::SequenceCache, ET_CHECK_OR_RETURN_ERROR( step.has_value(), InvalidArgument, "offgraph_kv: uncommittable step"); cache::SequenceCache::commit(*step); - metrics_.logical_length = length(); - DeviceGuard device(device_); - return mark_step_done(); + return pool_.mark_step_done(); } OffGraphKVMetrics metrics() const override { std::lock_guard guard(mutex_); - return metrics_; + OffGraphKVMetrics metrics; + metrics.logical_length = length(); + metrics.flat_capacity = pool_.rows(); + metrics.growth_count = pool_.growth_count(); + metrics.allocated_bytes = pool_.allocated_bytes(); + return metrics; } protected: @@ -263,23 +197,6 @@ class CudaSequenceKVCache final : public cache::SequenceCache, } private: - // Bytes in one sequence row of a layer: every head's vector at one position. - size_t row_bytes(const cache::LayerGeometry& layer) const { - return static_cast(layer.n_kv_heads) * - static_cast(layer.head_dim) * - slimc10::elementSize(storage_dtype_); - } - - // Slots a ring layer needs to serve one step of up to max_write tokens: the - // step writes all of them before attending, and its earliest query still - // reads back window - 1 positions. Same formula as cache::RingPolicy and as - // ring_physical_capacity() in backends/cuda/passes/lower_offgraph_kv.py. - // make_cuda_sequence_kv_cache() refuses ring layers without max_write. - int64_t ring_capacity(const cache::LayerGeometry& layer) const { - return static_cast(layer.policy.window) + - config_.max_write.value_or(1) - 1; - } - static Error check_write_length(int64_t write_length) { ET_CHECK_OR_RETURN_ERROR( write_length > 0 && write_length <= std::numeric_limits::max(), @@ -307,417 +224,15 @@ class CudaSequenceKVCache final : public cache::SequenceCache, return Error::Ok; } - // Rows the compiled program declares for a layer's storage. - int64_t declared_rows(const cache::LayerGeometry& layer) const { - return is_ring(layer) ? ring_capacity(layer) : config_.capacity; - } - - // Whole-model check, so it cannot run until every program has been loaded - // and registered. The first step is the earliest moment that is guaranteed. - // Caller holds mutex_. - Error validate_locked() const { - for (size_t index = 0; index < geometry_.layers.size(); ++index) { - for (const char* suffix : {"k", "v"}) { - if (discovered_fqns_.count(fqn(static_cast(index), suffix)) == - 0) { - ET_LOG( - Error, - "offgraph_kv: missing AOTI storage for layer %zu (%s)", - index, - suffix); - return Error::InvalidProgram; - } - } - } - return handles_associated_ ? Error::Ok : Error::InvalidState; - } - - Error allocate_layer( - const cache::LayerGeometry& layer, - int64_t rows, - cudaStream_t stream, - Allocation& out) { - const size_t bytes = row_bytes(layer) * static_cast(rows); - auto k = CudaAllocator::allocate_async(bytes, device_, stream); - ET_CHECK_OK_OR_RETURN_ERROR(k.error()); - auto v = CudaAllocator::allocate_async(bytes, device_, stream); - if (!v.ok()) { - CudaAllocator::deallocate_async(k.get(), device_, stream); - return v.error(); - } - out = Allocation{k.get(), v.get(), rows}; - metrics_.allocated_bytes += static_cast(2 * bytes); - return Error::Ok; - } - - void release(Allocation& allocation, cudaStream_t stream) { - if (!device_known_) { - return; - } - CudaAllocator::deallocate_async(allocation.k, device_, stream); - CudaAllocator::deallocate_async(allocation.v, device_, stream); - } - - // Releases an allocation made by allocate_layer and takes it off the books. - void discard( - const cache::LayerGeometry& layer, - Allocation& allocation, - cudaStream_t stream) { - metrics_.allocated_bytes -= - static_cast(2 * row_bytes(layer) * allocation.rows); - release(allocation, stream); - } - - // The previous step may have run on another stream (a caller-selected one), - // and its kernels may still be writing the storage this step reads, or that - // a growth below copies and frees. Order this stream behind it. - Error follow_previous_step(cudaStream_t stream) { - if (last_step_done_ == nullptr || stream == stream_) { - return Error::Ok; - } - const cudaError_t error = cudaStreamWaitEvent(stream, last_step_done_, 0); - ET_CHECK_OR_RETURN_ERROR( - error == cudaSuccess, - Internal, - "offgraph_kv: cannot order the step stream behind the previous one: %s", - cudaGetErrorString(error)); - return Error::Ok; - } - - // Records where the committed step's work ends on its stream. Called once - // the delegate has enqueued all of it, so the event covers every kernel that - // touched the storage. - Error mark_step_done() { - if (last_step_done_ == nullptr) { - const cudaError_t error = - cudaEventCreateWithFlags(&last_step_done_, cudaEventDisableTiming); - ET_CHECK_OR_RETURN_ERROR( - error == cudaSuccess, - Internal, - "offgraph_kv: cannot create the step event: %s", - cudaGetErrorString(error)); - } - const cudaError_t error = cudaEventRecord(last_step_done_, stream_); - ET_CHECK_OR_RETURN_ERROR( - error == cudaSuccess, - Internal, - "offgraph_kv: cannot record the step event: %s", - cudaGetErrorString(error)); - return Error::Ok; - } - - // A first step wider than the initial capacity is allocated at its own - // width rather than allocated and immediately grown. - Error ensure_initial_allocations(int64_t required, cudaStream_t stream) { - if (!allocations_.empty()) { - return Error::Ok; - } - const int64_t flat_rows = - std::max(config_.initial_capacity, required); - // All or nothing: a failure frees what was allocated so far, and the next - // step may try again. - std::vector allocated; - allocated.reserve(geometry_.layers.size()); - for (const cache::LayerGeometry& layer : geometry_.layers) { - const int64_t rows = is_ring(layer) ? ring_capacity(layer) : flat_rows; - Allocation allocation; - const Error error = allocate_layer(layer, rows, stream, allocation); - if (error != Error::Ok) { - for (size_t index = 0; index < allocated.size(); ++index) { - discard(geometry_.layers[index], allocated[index], stream); - } - return error; - } - allocated.push_back(allocation); - } - for (size_t index = 0; index < allocated.size(); ++index) { - allocations_.emplace(static_cast(index), allocated[index]); - } - metrics_.flat_capacity = flat_rows; - ET_LOG( - Info, - "offgraph_kv: initialized flat_capacity=%lld allocated_bytes=%lld", - static_cast(metrics_.flat_capacity), - static_cast(metrics_.allocated_bytes)); - return Error::Ok; - } - - // Reallocates every flat layer at new_rows and carries the rows already - // written across. BSHD storage makes that one contiguous prefix per buffer. - // Everything is ordered on `stream`, so the copy runs after the last step - // that wrote the old storage, the old storage is freed only after the copy, - // and the next step's kernels see the copied rows. - // - // Transactional: every replacement is allocated and filled before any old - // storage is released. A failure on any layer frees the replacements and - // leaves the cache exactly as it was -- storage, bindings, capacity -- so the - // step fails but the cache stays usable. - Error grow_flat(int64_t new_rows, cudaStream_t stream) { - const int64_t old_rows = metrics_.flat_capacity; - const int64_t live_rows = length(); - std::vector> replacements; - auto roll_back = [&]() { - for (auto& [index, replacement] : replacements) { - discard(geometry_.layers[index], replacement, stream); - } - }; - for (size_t index = 0; index < geometry_.layers.size(); ++index) { - const cache::LayerGeometry& layer = geometry_.layers[index]; - if (is_ring(layer)) { - continue; - } - Allocation replacement; - const Error error = allocate_layer(layer, new_rows, stream, replacement); - if (error != Error::Ok) { - roll_back(); - return error; - } - replacements.emplace_back(index, replacement); - } - for (const auto& [index, replacement] : replacements) { - const Allocation& current = allocations_.at(static_cast(index)); - const size_t live_bytes = - row_bytes(geometry_.layers[index]) * static_cast(live_rows); - if (live_bytes == 0) { - continue; - } - for (const auto& [dst, src] : - {std::pair{replacement.k, current.k}, - std::pair{replacement.v, current.v}}) { - const cudaError_t copy_error = cudaMemcpyAsync( - dst, src, live_bytes, cudaMemcpyDeviceToDevice, stream); - if (copy_error != cudaSuccess) { - ET_LOG( - Error, - "offgraph_kv: growth copy failed: %s", - cudaGetErrorString(copy_error)); - roll_back(); - return Error::Internal; - } - } - } - for (auto& [index, replacement] : replacements) { - Allocation& current = allocations_.at(static_cast(index)); - discard(geometry_.layers[index], current, stream); - current = replacement; - } - metrics_.flat_capacity = new_rows; - metrics_.growth_count++; - // Every program sharing this cache now points at freed storage: drop the - // bindings so each rebinds before its next run, and any captured CUDA - // graph so it is captured again against the new storage. prefill usually - // grows the cache while decode's graph sits idle, so this reaches every - // handle, not only the one stepping now. - // - // Rebinding also resets AOTI's constant-fold state, which must be run - // eagerly, so every graph-enabled handle gets at least one eager step - // before it captures -- including one that was about to capture for the - // first time, and without shortening a longer warmup still outstanding. - bound_.clear(); - for (auto& entry : descriptors_) { - CudaGraphState& graph = entry.first->cuda_graph_state; - if (graph.phase == CudaGraphPhase::Replay) { - graph.recapture(); - } else if (graph.phase == CudaGraphPhase::Warmup) { - graph.warmup_remaining = std::max(graph.warmup_remaining, 1); - } - } - ET_LOG( - Info, - "offgraph_kv: grew flat_capacity=%lld->%lld allocated_bytes=%lld " - "growth_count=%lld", - static_cast(old_rows), - static_cast(new_rows), - static_cast(metrics_.allocated_bytes), - static_cast(metrics_.growth_count)); - return Error::Ok; - } - - Error build_descriptors(CudaDelegateHandle* handle) { - ET_CHECK_OR_RETURN_ERROR( - handle->get_num_constants && handle->get_constant_name && - handle->get_constant_original_fqn && - handle->update_user_managed_constant_buffer_pairs, - NotSupported, - "offgraph_kv: AOTI external-buffer APIs are unavailable"); - size_t count = 0; - ET_CHECK_OK_OR_RETURN_ERROR( - handle->get_num_constants(handle->container_handle, &count)); - struct Compiled { - std::string internal_name; - size_t index; - }; - std::unordered_map compiled; - for (size_t index = 0; index < count; ++index) { - const char* internal = nullptr; - const char* original = nullptr; - ET_CHECK_OK_OR_RETURN_ERROR(handle->get_constant_name( - handle->container_handle, index, &internal)); - ET_CHECK_OK_OR_RETURN_ERROR(handle->get_constant_original_fqn( - handle->container_handle, index, &original)); - if (internal && original && internal[0] && original[0]) { - compiled.emplace(original, Compiled{internal, index}); - } - } - - // A reload of the same handle replaces what it had rather than adding to - // it, so its constants are never bound twice. - auto& descriptors = descriptors_[handle]; - descriptors.clear(); - bound_.erase(handle); - size_t found_layers = 0; - for (size_t index = 0; index < geometry_.layers.size(); ++index) { - const int64_t layer_id = static_cast(index); - size_t found = 0; - for (const auto& [suffix, is_value] : - {std::pair{"k", false}, std::pair{"v", true}}) { - const std::string name = fqn(layer_id, suffix); - const auto it = compiled.find(name); - if (it == compiled.end()) { - continue; - } - ET_CHECK_OK_OR_RETURN_ERROR(check_compiled( - handle, geometry_.layers[index], name, it->second.index)); - ++found; - descriptors.push_back( - Descriptor{it->second.internal_name, layer_id, is_value}); - discovered_fqns_.insert(name); - } - if (found == 2) { - ++found_layers; - } - } - // A method either has no off-graph storage (embeddings, vision) or has all - // of it. Anything between means the lowering pass and this runtime disagree - // about the geometry, which is worth failing on here -- while the offending - // method is still named -- rather than at the first decode. - ET_CHECK_OR_RETURN_ERROR( - found_layers == 0 || found_layers == geometry_.layers.size(), - InvalidProgram, - "offgraph_kv: program carries %zu of %zu layers' storage", - found_layers, - geometry_.layers.size()); - return Error::Ok; - } - - // The program's kernels address its storage with the shape and dtype it was - // compiled with, and AOTI binds external buffers without checking either. - // So this cache's own idea of a layer -- dtype, heads, head dim, and rows - // (capacity, or window + max_write - 1 for a ring) -- must match what the - // program declared, or a step could write past the allocation. - Error check_compiled( - CudaDelegateHandle* handle, - const cache::LayerGeometry& layer, - const std::string& name, - size_t index) const { - ET_CHECK_OR_RETURN_ERROR( - handle->get_constant_dtype && handle->get_constant_data_size, - NotSupported, - "offgraph_kv: AOTI constant metadata APIs are unavailable"); - int32_t dtype = 0; - size_t data_size = 0; - ET_CHECK_OK_OR_RETURN_ERROR( - handle->get_constant_dtype(handle->container_handle, index, &dtype)); - ET_CHECK_OK_OR_RETURN_ERROR(handle->get_constant_data_size( - handle->container_handle, index, &data_size)); - ET_CHECK_OR_RETURN_ERROR( - dtype == static_cast(storage_dtype_), - InvalidProgram, - "offgraph_kv: %s is compiled as dtype %d but the cache stores %d", - name.c_str(), - static_cast(dtype), - static_cast(storage_dtype_)); - const size_t expected = - row_bytes(layer) * static_cast(declared_rows(layer)); - ET_CHECK_OR_RETURN_ERROR( - data_size == expected, - InvalidProgram, - "offgraph_kv: %s is compiled with %zu bytes but the cache declares %zu; " - "the cache's capacity, max_write, window or head geometry does not " - "match the program", - name.c_str(), - data_size, - expected); - return Error::Ok; - } - - // Binds each storage constant with the shape the program declared (BSHD at - // the maximum rows) over the current allocation, which may hold fewer rows. - // That is safe because every access the program makes is bounded by kv_len - // along the sequence, and prepare_step() has grown the allocation past it. - Error bind(CudaDelegateHandle* handle) { - if (bound_.find(handle) != bound_.end()) { - return Error::Ok; - } - const std::vector& descriptors = descriptors_[handle]; - Bound bound; - bound.tensors.reserve(descriptors.size()); - std::vector pairs; - pairs.reserve(descriptors.size()); - for (const Descriptor& descriptor : descriptors) { - const cache::LayerGeometry& layer = - geometry_.layers[static_cast(descriptor.layer_id)]; - const Allocation& allocation = allocations_.at(descriptor.layer_id); - const int64_t declared = declared_rows(layer); - ET_CHECK_OR_RETURN_ERROR( - allocation.rows <= declared, - Internal, - "offgraph_kv: layer %lld holds %lld rows, more than the %lld declared", - static_cast(descriptor.layer_id), - static_cast(allocation.rows), - static_cast(declared)); - const int64_t heads = layer.n_kv_heads; - const int64_t dim = layer.head_dim; - const int64_t sizes[] = {1, declared, heads, dim}; - const int64_t strides[] = {declared * heads * dim, heads * dim, dim, 1}; - auto tensor = std::make_unique(from_blob( - descriptor.is_value ? allocation.v : allocation.k, - ::executorch::runtime::makeArrayRef(sizes, 4), - ::executorch::runtime::makeArrayRef(strides, 4), - storage_dtype_, - slimc10::Device(slimc10::DeviceType::CUDA, device_))); - pairs.push_back( - {descriptor.internal_name.c_str(), - reinterpret_cast(tensor.get())}); - bound.tensors.push_back(std::move(tensor)); - } - if (!pairs.empty()) { - ET_CHECK_OK_OR_RETURN_ERROR( - handle->update_user_managed_constant_buffer_pairs( - handle->container_handle, - pairs.data(), - pairs.size(), - /*use_inactive=*/false, - /*validate_full_update=*/false)); - } - bound_.emplace(handle, std::move(bound)); - return Error::Ok; - } - // Guards everything below. The engine serialises its own calls, but the // delegate reaches note_handle/rebind from whichever thread loads or runs a // method. mutable std::mutex mutex_; cache::CacheGeometry geometry_; - cache::CacheConfig config_; - slimc10::ScalarType storage_dtype_; - - int device_{0}; - bool device_known_{false}; - bool handles_associated_{false}; + CudaKVPool pool_; bool validated_{false}; - // The stream of the latest step, and so of the latest use of the storage. - cudaStream_t stream_{cudaStreamPerThread}; - // Recorded on stream_ when a step commits; a later step on another stream - // waits on it before touching the storage. - cudaEvent_t last_step_done_{nullptr}; Error error_{Error::Ok}; - OffGraphKVMetrics metrics_; - std::unordered_map allocations_; - std::unordered_map> descriptors_; - std::unordered_map bound_; - std::unordered_set discovered_fqns_; }; } // namespace diff --git a/backends/cuda/runtime/cuda_kv_pool.cpp b/backends/cuda/runtime/cuda_kv_pool.cpp new file mode 100644 index 00000000000..d23dc54de4d --- /dev/null +++ b/backends/cuda/runtime/cuda_kv_pool.cpp @@ -0,0 +1,637 @@ +/* + * Copyright (c) Meta Platforms, Inc. and affiliates. + * All rights reserved. + * + * This source code is licensed under the BSD-style license found in the + * LICENSE file in the root directory of this source tree. + */ + +#include + +#include +#include +#include + +#include +#include +#include +#include + +namespace executorch::backends::cuda { + +namespace aoti = ::executorch::backends::aoti; +namespace slimc10 = ::executorch::backends::aoti::slim::c10; +using ::executorch::backends::aoti::slim::from_blob; +using ::executorch::backends::aoti::slim::SlimTensor; +using ::executorch::runtime::Error; + +std::string offgraph_kv_layer_fqn(int64_t layer_id, const char* suffix) { + return "__et_offgraph_kv_layer_" + std::to_string(layer_id) + "_" + suffix; +} + +namespace { + +int64_t numel(const std::vector& sizes) { + int64_t n = 1; + for (const int64_t size : sizes) { + n *= size; + } + return n; +} + +std::vector contiguous_strides(const std::vector& sizes) { + std::vector strides(sizes.size(), 1); + for (size_t i = sizes.size(); i > 1; --i) { + strides[i - 2] = strides[i - 1] * sizes[i - 1]; + } + return strides; +} + +// Makes the pool's device current for a scope. The backend runs on one +// device, so this is normally a no-op, but event and stream calls act on the +// current device and a caller may have switched it. +class DeviceGuard { + public: + explicit DeviceGuard(int device) { + if (cudaGetDevice(&previous_) == cudaSuccess && previous_ != device) { + restore_ = cudaSetDevice(device) == cudaSuccess; + } + } + ~DeviceGuard() { + if (restore_) { + (void)cudaSetDevice(previous_); + } + } + DeviceGuard(const DeviceGuard&) = delete; + DeviceGuard& operator=(const DeviceGuard&) = delete; + + private: + int previous_{0}; + bool restore_{false}; +}; + +} // namespace + +CudaKVPool::CudaKVPool( + std::vector layers, + std::vector side_buffers, + slimc10::ScalarType storage_dtype, + int64_t initial_rows) + : layers_(std::move(layers)), + side_specs_(std::move(side_buffers)), + storage_dtype_(storage_dtype), + initial_rows_(initial_rows) { + max_rows_ = std::numeric_limits::max(); + for (const Layer& layer : layers_) { + if (layer.growable) { + max_rows_ = std::min(max_rows_, layer.declared_rows); + } + } +} + +// Nothing else can reach the pool once it is being destroyed. It never +// touches stream_: that is the caller's stream, which may already be gone. +// The device is drained instead, which covers any step still in flight on +// whatever stream ran it. +CudaKVPool::~CudaKVPool() { + if (!device_known_) { + return; + } + DeviceGuard device(device_); + (void)cudaDeviceSynchronize(); + for (auto& allocation : allocations_) { + release(allocation.k, cudaStreamPerThread); + release(allocation.v, cudaStreamPerThread); + } + for (void* buffer : side_buffers_) { + release(buffer, cudaStreamPerThread); + } + (void)cudaStreamSynchronize(cudaStreamPerThread); + if (last_step_done_ != nullptr) { + (void)cudaEventDestroy(last_step_done_); + } +} + +runtime::Result CudaKVPool::note_handle(CudaDelegateHandle* handle) { + if (!device_known_) { + ET_CHECK_OR_RETURN_ERROR( + cudaGetDevice(&device_) == cudaSuccess, + Internal, + "offgraph_kv: cannot query the current CUDA device"); + device_known_ = true; + } + ET_CHECK_OK_OR_RETURN_ERROR(build_descriptors(handle)); + if (descriptors_[handle].empty()) { + descriptors_.erase(handle); + return false; + } + return true; +} + +void CudaKVPool::forget_handle(CudaDelegateHandle* handle) { + descriptors_.erase(handle); + bound_.erase(handle); +} + +bool CudaKVPool::serves(CudaDelegateHandle* handle) const { + return descriptors_.find(handle) != descriptors_.end(); +} + +Error CudaKVPool::validate() const { + for (size_t index = 0; index < layers_.size(); ++index) { + for (const char* suffix : {"k", "v"}) { + if (discovered_fqns_.count( + offgraph_kv_layer_fqn(static_cast(index), suffix)) == + 0) { + ET_LOG( + Error, + "offgraph_kv: missing AOTI storage for layer %zu (%s)", + index, + suffix); + return Error::InvalidProgram; + } + } + } + for (const SideBuffer& spec : side_specs_) { + if (discovered_fqns_.count(spec.fqn) == 0) { + ET_LOG( + Error, "offgraph_kv: missing AOTI side buffer %s", spec.fqn.c_str()); + return Error::InvalidProgram; + } + } + return Error::Ok; +} + +Error CudaKVPool::prepare( + int64_t required_rows, + int64_t live_rows, + cudaStream_t stream) { + DeviceGuard device(device_); + ET_CHECK_OK_OR_RETURN_ERROR(follow_previous_step(stream)); + stream_ = stream; + if (!allocated_) { + ET_CHECK_OK_OR_RETURN_ERROR(allocate_initial(required_rows, stream)); + } + if (required_rows > rows_) { + const int64_t next = + std::min(max_rows_, std::max(required_rows, rows_ * 2)); + ET_CHECK_OK_OR_RETURN_ERROR(grow(next, live_rows, stream)); + } + return Error::Ok; +} + +Error CudaKVPool::mark_step_done() { + DeviceGuard device(device_); + if (last_step_done_ == nullptr) { + const cudaError_t error = + cudaEventCreateWithFlags(&last_step_done_, cudaEventDisableTiming); + ET_CHECK_OR_RETURN_ERROR( + error == cudaSuccess, + Internal, + "offgraph_kv: cannot create the step event: %s", + cudaGetErrorString(error)); + } + const cudaError_t error = cudaEventRecord(last_step_done_, stream_); + ET_CHECK_OR_RETURN_ERROR( + error == cudaSuccess, + Internal, + "offgraph_kv: cannot record the step event: %s", + cudaGetErrorString(error)); + return Error::Ok; +} + +void* CudaKVPool::side_buffer(size_t index) const { + return index < side_buffers_.size() ? side_buffers_[index] : nullptr; +} + +size_t CudaKVPool::side_buffer_bytes(size_t index) const { + const SideBuffer& spec = side_specs_.at(index); + return static_cast(numel(spec.sizes)) * slimc10::elementSize(spec.dtype); +} + +size_t CudaKVPool::row_bytes(const Layer& layer) const { + return static_cast(layer.n_kv_heads) * + static_cast(layer.head_dim) * slimc10::elementSize(storage_dtype_); +} + +Error CudaKVPool::allocate_layer( + const Layer& layer, + int64_t rows, + cudaStream_t stream, + Allocation& out) { + const size_t bytes = row_bytes(layer) * static_cast(rows); + auto k = CudaAllocator::allocate_async(bytes, device_, stream); + ET_CHECK_OK_OR_RETURN_ERROR(k.error()); + auto v = CudaAllocator::allocate_async(bytes, device_, stream); + if (!v.ok()) { + CudaAllocator::deallocate_async(k.get(), device_, stream); + return v.error(); + } + out = Allocation{k.get(), v.get(), rows}; + allocated_bytes_ += static_cast(2 * bytes); + return Error::Ok; +} + +void CudaKVPool::release(void* ptr, cudaStream_t stream) { + if (device_known_ && ptr != nullptr) { + CudaAllocator::deallocate_async(ptr, device_, stream); + } +} + +void CudaKVPool::discard( + const Layer& layer, + Allocation& allocation, + cudaStream_t stream) { + allocated_bytes_ -= static_cast(2 * row_bytes(layer) * allocation.rows); + release(allocation.k, stream); + release(allocation.v, stream); + allocation = Allocation{}; +} + +// The previous step may have run on another stream (a caller-selected one), +// and its kernels may still be writing the storage this step reads, or that a +// growth below copies and frees. Order this stream behind it. +Error CudaKVPool::follow_previous_step(cudaStream_t stream) { + if (last_step_done_ == nullptr || stream == stream_) { + return Error::Ok; + } + const cudaError_t error = cudaStreamWaitEvent(stream, last_step_done_, 0); + ET_CHECK_OR_RETURN_ERROR( + error == cudaSuccess, + Internal, + "offgraph_kv: cannot order the step stream behind the previous one: %s", + cudaGetErrorString(error)); + return Error::Ok; +} + +// A first step wider than the initial rows is allocated at its own width +// rather than allocated and immediately grown. All or nothing: a failure frees +// what was allocated so far, and the next step may try again. +Error CudaKVPool::allocate_initial(int64_t required_rows, cudaStream_t stream) { + const int64_t growable_rows = std::max(initial_rows_, required_rows); + std::vector allocated; + allocated.reserve(layers_.size()); + for (const Layer& layer : layers_) { + const int64_t rows = layer.growable ? growable_rows : layer.declared_rows; + Allocation allocation; + const Error error = allocate_layer(layer, rows, stream, allocation); + if (error != Error::Ok) { + for (size_t index = 0; index < allocated.size(); ++index) { + discard(layers_[index], allocated[index], stream); + } + return error; + } + allocated.push_back(allocation); + } + const Error side_error = allocate_side_buffers(stream); + if (side_error != Error::Ok) { + for (size_t index = 0; index < allocated.size(); ++index) { + discard(layers_[index], allocated[index], stream); + } + return side_error; + } + allocations_ = std::move(allocated); + rows_ = growable_rows; + allocated_ = true; + ET_LOG( + Info, + "offgraph_kv: initialized flat_capacity=%lld allocated_bytes=%lld", + static_cast(rows_), + static_cast(allocated_bytes_)); + return Error::Ok; +} + +// Side buffers start zeroed, so a program bound before its first write reads +// a well-defined value rather than whatever the allocator handed back. +Error CudaKVPool::allocate_side_buffers(cudaStream_t stream) { + std::vector buffers; + buffers.reserve(side_specs_.size()); + auto roll_back = [&]() { + for (void* buffer : buffers) { + release(buffer, stream); + } + }; + for (size_t index = 0; index < side_specs_.size(); ++index) { + const size_t bytes = side_buffer_bytes(index); + auto buffer = CudaAllocator::allocate_async(bytes, device_, stream); + if (!buffer.ok()) { + roll_back(); + return buffer.error(); + } + buffers.push_back(buffer.get()); + const cudaError_t error = cudaMemsetAsync(buffer.get(), 0, bytes, stream); + if (error != cudaSuccess) { + roll_back(); + ET_LOG( + Error, + "offgraph_kv: cannot clear side buffer %s: %s", + side_specs_[index].fqn.c_str(), + cudaGetErrorString(error)); + return Error::Internal; + } + } + side_buffers_ = std::move(buffers); + return Error::Ok; +} + +// Reallocates every growable layer at new_rows and carries the rows already +// written across. BSHD storage makes that one contiguous prefix per buffer. +// Everything is ordered on `stream`, so the copy runs after the last step +// that wrote the old storage, the old storage is freed only after the copy, +// and the next step's kernels see the copied rows. +// +// Transactional: every replacement is allocated and filled before any old +// storage is released. A failure on any layer frees the replacements and +// leaves the pool exactly as it was -- storage, bindings, rows -- so the step +// fails but the pool stays usable. +Error CudaKVPool::grow(int64_t new_rows, int64_t live_rows, cudaStream_t stream) { + const int64_t old_rows = rows_; + std::vector> replacements; + auto roll_back = [&]() { + for (auto& [index, replacement] : replacements) { + discard(layers_[index], replacement, stream); + } + }; + for (size_t index = 0; index < layers_.size(); ++index) { + const Layer& layer = layers_[index]; + if (!layer.growable) { + continue; + } + Allocation replacement; + const Error error = allocate_layer(layer, new_rows, stream, replacement); + if (error != Error::Ok) { + roll_back(); + return error; + } + replacements.emplace_back(index, replacement); + } + for (const auto& [index, replacement] : replacements) { + const Allocation& current = allocations_.at(index); + const size_t live_bytes = + row_bytes(layers_[index]) * static_cast(live_rows); + if (live_bytes == 0) { + continue; + } + for (const auto& [dst, src] : + {std::pair{replacement.k, current.k}, + std::pair{replacement.v, current.v}}) { + const cudaError_t copy_error = cudaMemcpyAsync( + dst, src, live_bytes, cudaMemcpyDeviceToDevice, stream); + if (copy_error != cudaSuccess) { + ET_LOG( + Error, + "offgraph_kv: growth copy failed: %s", + cudaGetErrorString(copy_error)); + roll_back(); + return Error::Internal; + } + } + } + for (auto& [index, replacement] : replacements) { + Allocation& current = allocations_.at(index); + discard(layers_[index], current, stream); + current = replacement; + } + rows_ = new_rows; + growth_count_++; + // Every program sharing this pool now points at freed storage: drop the + // bindings so each rebinds before its next run, and any captured CUDA graph + // so it is captured again against the new storage. Prefill usually grows + // the cache while decode's graph sits idle, so this reaches every handle, + // not only the one stepping now. + // + // Rebinding also resets AOTI's constant-fold state, which must be run + // eagerly, so every graph-enabled handle gets at least one eager step + // before it captures -- including one that was about to capture for the + // first time, and without shortening a longer warmup still outstanding. + bound_.clear(); + for (auto& entry : descriptors_) { + CudaGraphState& graph = entry.first->cuda_graph_state; + if (graph.phase == CudaGraphPhase::Replay) { + graph.recapture(); + } else if (graph.phase == CudaGraphPhase::Warmup) { + graph.warmup_remaining = std::max(graph.warmup_remaining, 1); + } + } + ET_LOG( + Info, + "offgraph_kv: grew flat_capacity=%lld->%lld allocated_bytes=%lld " + "growth_count=%lld", + static_cast(old_rows), + static_cast(new_rows), + static_cast(allocated_bytes_), + static_cast(growth_count_)); + return Error::Ok; +} + +Error CudaKVPool::build_descriptors(CudaDelegateHandle* handle) { + ET_CHECK_OR_RETURN_ERROR( + handle->get_num_constants && handle->get_constant_name && + handle->get_constant_original_fqn && + handle->update_user_managed_constant_buffer_pairs, + NotSupported, + "offgraph_kv: AOTI external-buffer APIs are unavailable"); + size_t count = 0; + ET_CHECK_OK_OR_RETURN_ERROR( + handle->get_num_constants(handle->container_handle, &count)); + struct Compiled { + std::string internal_name; + size_t index; + }; + std::unordered_map compiled; + for (size_t index = 0; index < count; ++index) { + const char* internal = nullptr; + const char* original = nullptr; + ET_CHECK_OK_OR_RETURN_ERROR( + handle->get_constant_name(handle->container_handle, index, &internal)); + ET_CHECK_OK_OR_RETURN_ERROR(handle->get_constant_original_fqn( + handle->container_handle, index, &original)); + if (internal && original && internal[0] && original[0]) { + compiled.emplace(original, Compiled{internal, index}); + } + } + + std::vector descriptors; + std::vector found_fqns; + size_t found_layers = 0; + for (size_t index = 0; index < layers_.size(); ++index) { + size_t found = 0; + for (const auto& [suffix, slot] : + {std::pair{"k", Slot::Key}, std::pair{"v", Slot::Value}}) { + const std::string name = + offgraph_kv_layer_fqn(static_cast(index), suffix); + const auto it = compiled.find(name); + if (it == compiled.end()) { + continue; + } + const Layer& layer = layers_[index]; + ET_CHECK_OK_OR_RETURN_ERROR(check_compiled( + handle, + it->second.index, + name, + storage_dtype_, + row_bytes(layer) * static_cast(layer.declared_rows))); + ++found; + descriptors.push_back(Descriptor{it->second.internal_name, slot, index}); + found_fqns.push_back(name); + } + if (found == 2) { + ++found_layers; + } + } + size_t found_side = 0; + for (size_t index = 0; index < side_specs_.size(); ++index) { + const SideBuffer& spec = side_specs_[index]; + const auto it = compiled.find(spec.fqn); + if (it != compiled.end()) { + ET_CHECK_OK_OR_RETURN_ERROR(check_compiled( + handle, + it->second.index, + spec.fqn, + spec.dtype, + side_buffer_bytes(index))); + ++found_side; + descriptors.push_back( + Descriptor{it->second.internal_name, Slot::Side, index}); + found_fqns.push_back(spec.fqn); + } + } + // A method either has no off-graph storage (embeddings, vision) or has all + // of it. Anything between means the lowering pass and this runtime disagree + // about the geometry, which is worth failing on here -- while the offending + // method is still named -- rather than at the first step. + ET_CHECK_OR_RETURN_ERROR( + found_layers == 0 || found_layers == layers_.size(), + InvalidProgram, + "offgraph_kv: program carries %zu of %zu layers' storage", + found_layers, + layers_.size()); + ET_CHECK_OR_RETURN_ERROR( + found_side == (found_layers == 0 ? 0 : side_specs_.size()), + InvalidProgram, + "offgraph_kv: program carries %zu of %zu side buffers", + found_side, + side_specs_.size()); + discovered_fqns_.insert(found_fqns.begin(), found_fqns.end()); + // A reload of the same handle replaces what it had rather than adding to + // it, so its constants are never bound twice. + descriptors_[handle] = std::move(descriptors); + bound_.erase(handle); + return Error::Ok; +} + +// The program's kernels address its storage with the shape and dtype it was +// compiled with, and AOTI binds external buffers without checking either. So +// the pool's idea of each constant -- dtype, and bytes (heads x head dim x +// declared rows for a layer, the declared shape for a side buffer) -- must +// match what the program declared, or a step could write past the allocation. +// +// AOTI reports a constant's storage bytes, rounded up to a multiple of 64 +// when the program also holds CPU constants (cpp_wrapper_cpu.py), so either +// form matches; the rounding never exceeds 63 bytes and never shrinks. +Error CudaKVPool::check_compiled( + CudaDelegateHandle* handle, + size_t constant_index, + const std::string& name, + slimc10::ScalarType dtype, + size_t bytes) const { + ET_CHECK_OR_RETURN_ERROR( + handle->get_constant_dtype && handle->get_constant_data_size, + NotSupported, + "offgraph_kv: AOTI constant metadata APIs are unavailable"); + int32_t compiled_dtype = 0; + size_t compiled_bytes = 0; + ET_CHECK_OK_OR_RETURN_ERROR(handle->get_constant_dtype( + handle->container_handle, constant_index, &compiled_dtype)); + ET_CHECK_OK_OR_RETURN_ERROR(handle->get_constant_data_size( + handle->container_handle, constant_index, &compiled_bytes)); + ET_CHECK_OR_RETURN_ERROR( + compiled_dtype == static_cast(dtype), + InvalidProgram, + "offgraph_kv: %s is compiled as dtype %d but the cache stores %d", + name.c_str(), + static_cast(compiled_dtype), + static_cast(dtype)); + constexpr size_t kAotiConstantAlignment = 64; + const size_t aligned = + (bytes + kAotiConstantAlignment - 1) / kAotiConstantAlignment * + kAotiConstantAlignment; + ET_CHECK_OR_RETURN_ERROR( + compiled_bytes == bytes || compiled_bytes == aligned, + InvalidProgram, + "offgraph_kv: %s is compiled with %zu bytes but the cache declares %zu; " + "the cache's capacity, max_write, window or head geometry does not " + "match the program", + name.c_str(), + compiled_bytes, + bytes); + return Error::Ok; +} + +// Binds each constant with the shape the program declared over the current +// allocation. A growable layer is declared BSHD at its maximum rows but may be +// backed by fewer: every access the program makes is bounded by kv_len along +// the sequence, and prepare() has grown the allocation past it. +Error CudaKVPool::bind(CudaDelegateHandle* handle) { + if (bound_.find(handle) != bound_.end()) { + return Error::Ok; + } + const auto descriptors = descriptors_.find(handle); + if (descriptors == descriptors_.end()) { + return Error::Ok; + } + ET_CHECK_OR_RETURN_ERROR( + allocated_, InvalidState, "offgraph_kv: prepare must run before bind"); + Bound bound; + bound.tensors.reserve(descriptors->second.size()); + std::vector pairs; + pairs.reserve(descriptors->second.size()); + const slimc10::Device device(slimc10::DeviceType::CUDA, device_); + for (const Descriptor& descriptor : descriptors->second) { + void* data = nullptr; + std::vector sizes; + slimc10::ScalarType dtype = storage_dtype_; + if (descriptor.slot == Slot::Side) { + const SideBuffer& spec = side_specs_[descriptor.index]; + data = side_buffers_[descriptor.index]; + sizes = spec.sizes; + dtype = spec.dtype; + } else { + const Layer& layer = layers_[descriptor.index]; + const Allocation& allocation = allocations_[descriptor.index]; + ET_CHECK_OR_RETURN_ERROR( + allocation.rows <= layer.declared_rows, + Internal, + "offgraph_kv: layer %zu holds %lld rows, more than the %lld declared", + descriptor.index, + static_cast(allocation.rows), + static_cast(layer.declared_rows)); + data = descriptor.slot == Slot::Value ? allocation.v : allocation.k; + sizes = {1, layer.declared_rows, layer.n_kv_heads, layer.head_dim}; + } + const std::vector strides = contiguous_strides(sizes); + auto tensor = std::make_unique(from_blob( + data, + ::executorch::runtime::makeArrayRef(sizes.data(), sizes.size()), + ::executorch::runtime::makeArrayRef(strides.data(), strides.size()), + dtype, + device)); + pairs.push_back( + {descriptor.internal_name.c_str(), + reinterpret_cast(tensor.get())}); + bound.tensors.push_back(std::move(tensor)); + } + if (!pairs.empty()) { + ET_CHECK_OK_OR_RETURN_ERROR( + handle->update_user_managed_constant_buffer_pairs( + handle->container_handle, + pairs.data(), + pairs.size(), + /*use_inactive=*/false, + /*validate_full_update=*/false)); + } + bound_.emplace(handle, std::move(bound)); + return Error::Ok; +} + +} // namespace executorch::backends::cuda diff --git a/backends/cuda/runtime/cuda_kv_pool.h b/backends/cuda/runtime/cuda_kv_pool.h new file mode 100644 index 00000000000..46d8c0f44ce --- /dev/null +++ b/backends/cuda/runtime/cuda_kv_pool.h @@ -0,0 +1,184 @@ +/* + * Copyright (c) Meta Platforms, Inc. and affiliates. + * All rights reserved. + * + * This source code is licensed under the BSD-style license found in the + * LICENSE file in the root directory of this source tree. + */ + +#pragma once + +#include +#include +#include +#include +#include +#include + +#include +#include +#include +#include +#include +#include + +namespace executorch::backends::cuda { + +// The bytes behind an off-graph KV cache on one CUDA device: per-layer K/V row +// storage, plus fixed-size side buffers a layout may declare alongside it, and +// the AOTI bindings that point every compiled program at them. Which rows mean +// what is the owning cache's business, never this class's. +// +// Growable layers share one row count that grows geometrically toward their +// declared rows, moving their storage; fixed layers and side buffers are +// allocated once and never move, so a program may capture their addresses. +// +// Not thread-safe: the owning cache serializes every call. +class CudaKVPool final { + public: + struct Layer { + int64_t n_kv_heads{0}; + int64_t head_dim{0}; + // Rows the compiled program declares for this layer's storage. + int64_t declared_rows{0}; + bool growable{false}; + }; + + // A constant the program declares by FQN with a fixed, contiguous shape. + struct SideBuffer { + std::string fqn; + ::executorch::backends::aoti::slim::c10::ScalarType dtype{ + ::executorch::backends::aoti::slim::c10::ScalarType::Byte}; + std::vector sizes; + }; + + CudaKVPool( + std::vector layers, + std::vector side_buffers, + ::executorch::backends::aoti::slim::c10::ScalarType storage_dtype, + int64_t initial_rows); + + CudaKVPool(const CudaKVPool&) = delete; + CudaKVPool& operator=(const CudaKVPool&) = delete; + CudaKVPool(CudaKVPool&&) = delete; + CudaKVPool& operator=(CudaKVPool&&) = delete; + + ~CudaKVPool(); + + // Load time. Discovers this program's storage constants; false = it carries + // none. A program carrying some layers but not all, or layers without every + // side buffer, is rejected as disagreeing with the pool's geometry. + runtime::Result note_handle(CudaDelegateHandle* handle); + void forget_handle(CudaDelegateHandle* handle); + bool serves(CudaDelegateHandle* handle) const; + + // Whole-model check: every layer and side buffer was found in some program. + runtime::Error validate() const; + + // Orders `stream` after the previous step, then makes the growable layers + // hold at least `required_rows`, carrying their first `live_rows` rows + // across a growth. The first call also allocates the fixed storage. + runtime::Error + prepare(int64_t required_rows, int64_t live_rows, cudaStream_t stream); + + // Records where the step's work ends on its stream, for the next step to + // follow if it runs on another. + runtime::Error mark_step_done(); + + // Points `handle`'s constants at the current storage. Cached until storage + // moves. Precondition: allocated(). + runtime::Error bind(CudaDelegateHandle* handle); + + bool allocated() const { + return allocated_; + } + // Device pointer of side buffer `index`, fixed for the pool's life once + // allocated(). + void* side_buffer(size_t index) const; + size_t side_buffer_bytes(size_t index) const; + // The stream of the latest step, on which side-buffer writes must be issued. + cudaStream_t stream() const { + return stream_; + } + + int64_t rows() const { + return rows_; + } + int64_t allocated_bytes() const { + return allocated_bytes_; + } + int64_t growth_count() const { + return growth_count_; + } + + private: + struct Allocation { + void* k{nullptr}; + void* v{nullptr}; + int64_t rows{0}; + }; + + enum class Slot { Key, Value, Side }; + + struct Descriptor { + std::string internal_name; + Slot slot{Slot::Key}; + size_t index{0}; + }; + + // The tensors a handle's AOTI constants currently point at. + struct Bound { + std::vector> + tensors; + }; + + size_t row_bytes(const Layer& layer) const; + runtime::Error allocate_layer( + const Layer& layer, + int64_t rows, + cudaStream_t stream, + Allocation& out); + void release(void* ptr, cudaStream_t stream); + void discard(const Layer& layer, Allocation& allocation, cudaStream_t stream); + runtime::Error follow_previous_step(cudaStream_t stream); + runtime::Error allocate_initial(int64_t required_rows, cudaStream_t stream); + runtime::Error allocate_side_buffers(cudaStream_t stream); + runtime::Error grow(int64_t new_rows, int64_t live_rows, cudaStream_t stream); + runtime::Error build_descriptors(CudaDelegateHandle* handle); + runtime::Error check_compiled( + CudaDelegateHandle* handle, + size_t constant_index, + const std::string& name, + ::executorch::backends::aoti::slim::c10::ScalarType dtype, + size_t bytes) const; + + std::vector layers_; + std::vector side_specs_; + ::executorch::backends::aoti::slim::c10::ScalarType storage_dtype_; + int64_t initial_rows_; + int64_t max_rows_{0}; + + int device_{0}; + bool device_known_{false}; + bool allocated_{false}; + // The stream of the latest step, and so of the latest use of the storage. + cudaStream_t stream_{cudaStreamPerThread}; + // Recorded on stream_ when a step commits; a later step on another stream + // waits on it before touching the storage. + cudaEvent_t last_step_done_{nullptr}; + + int64_t rows_{0}; + int64_t allocated_bytes_{0}; + int64_t growth_count_{0}; + std::vector allocations_; + std::vector side_buffers_; + + std::unordered_map> descriptors_; + std::unordered_map bound_; + std::unordered_set discovered_fqns_; +}; + +// The FQN the lowering pass gives a layer's K or V storage constant. +std::string offgraph_kv_layer_fqn(int64_t layer_id, const char* suffix); + +} // namespace executorch::backends::cuda diff --git a/backends/cuda/runtime/targets.bzl b/backends/cuda/runtime/targets.bzl index d4bb81d310e..2053b6e35c4 100644 --- a/backends/cuda/runtime/targets.bzl +++ b/backends/cuda/runtime/targets.bzl @@ -124,6 +124,7 @@ def define_common_targets(is_fbcode = False): srcs = [ "cuda_backend.cpp", "cuda_kv_cache.cpp", + "cuda_kv_pool.cpp", "cuda_mutable_state.cpp", "cuda_weight_cache.cpp", ], @@ -131,6 +132,7 @@ def define_common_targets(is_fbcode = False): "backend_options.h", "cuda_delegate_handle.h", "cuda_kv_cache.h", + "cuda_kv_pool.h", "cuda_mutable_state.h", "cuda_weight_cache.h", ], diff --git a/backends/cuda/runtime/test/test_cuda_kv_cache.cpp b/backends/cuda/runtime/test/test_cuda_kv_cache.cpp index 264a2a8de66..f78c801705f 100644 --- a/backends/cuda/runtime/test/test_cuda_kv_cache.cpp +++ b/backends/cuda/runtime/test/test_cuda_kv_cache.cpp @@ -8,12 +8,14 @@ #include #include +#include #include #include #include #include #include +#include #include #include #include @@ -46,8 +48,24 @@ struct FakeContainer { std::unordered_map data_size_override{}; std::unordered_map dtype_override{}; size_t last_update_pairs{0}; + // Compiled dtype and bytes by FQN, for constants make_cache() cannot + // describe: storage built straight on a CudaKVPool, and side buffers. + std::unordered_map> declared{}; }; +// Declares `fqn` as compiled with `dtype` and contiguous `sizes`. +void declare( + FakeContainer& container, + const std::string& fqn, + slimc10::ScalarType dtype, + std::initializer_list sizes) { + size_t bytes = slimc10::elementSize(dtype); + for (const int64_t size : sizes) { + bytes *= static_cast(size); + } + container.declared[fqn] = {static_cast(dtype), bytes}; +} + // What the program under test was "compiled" with, for the fake's defaults. struct Compiled { cache::CacheGeometry geometry; @@ -113,7 +131,12 @@ Error get_constant_dtype( int32_t* dtype) { auto* fake = reinterpret_cast(container); const auto it = fake->dtype_override.find(index); - *dtype = it != fake->dtype_override.end() ? it->second + if (it != fake->dtype_override.end()) { + *dtype = it->second; + return Error::Ok; + } + const auto declared = fake->declared.find(fake->fqns.at(index)); + *dtype = declared != fake->declared.end() ? declared->second.first : compiled().config.kv_dtype; return Error::Ok; } @@ -129,6 +152,11 @@ Error get_constant_data_size( return Error::Ok; } const std::string& fqn = fake->fqns.at(index); + const auto declared = fake->declared.find(fqn); + if (declared != fake->declared.end()) { + *data_size = declared->second.second; + return Error::Ok; + } *data_size = declared_bytes( compiled().geometry.layers.at(layer_of(fqn)), compiled().config); return Error::Ok; @@ -1005,3 +1033,265 @@ TEST_F(CudaKVCacheTest, GrowthDuringAnotherMethodDropsItsGraph) { kv.forget_handle(&prefill); kv.forget_handle(&decode); } + +namespace { + +// One growable layer, one fixed layer and two side buffers: everything the +// pool distinguishes, and nothing about what the rows mean. +class CudaKVPoolTest : public CudaKVCacheTest { + protected: + static constexpr int kDeclaredRows = 64; + static constexpr int kFixedRows = 6; + + std::unique_ptr make_pool(int initial_rows) { + return std::make_unique( + std::vector{ + {kHeads, kDim, kDeclaredRows, /*growable=*/true}, + {kHeads, kDim, kFixedRows, /*growable=*/false}, + }, + std::vector{ + {"__et_offgraph_kv_cells", slimc10::ScalarType::Long, {16}}, + {"__et_offgraph_kv_mask_w0", + slimc10::ScalarType::Bool, + {1, 1, 16, kDeclaredRows}}, + }, + slimc10::ScalarType::BFloat16, + initial_rows); + } + + // What a program lowered for make_pool() declares, for `fqns`. + static void declare_pool(FakeContainer& container) { + const auto bf16 = slimc10::ScalarType::BFloat16; + for (const char* fqn : + {"__et_offgraph_kv_layer_0_k", "__et_offgraph_kv_layer_0_v"}) { + declare(container, fqn, bf16, {1, kDeclaredRows, kHeads, kDim}); + } + for (const char* fqn : + {"__et_offgraph_kv_layer_1_k", "__et_offgraph_kv_layer_1_v"}) { + declare(container, fqn, bf16, {1, kFixedRows, kHeads, kDim}); + } + declare(container, "__et_offgraph_kv_cells", slimc10::ScalarType::Long, {16}); + declare( + container, + "__et_offgraph_kv_mask_w0", + slimc10::ScalarType::Bool, + {1, 1, 16, kDeclaredRows}); + } + + static FakeContainer full_container() { + FakeContainer container{ + {"grow_k", "grow_v", "fixed_k", "fixed_v", "cells", "mask"}, + {"__et_offgraph_kv_layer_0_k", + "__et_offgraph_kv_layer_0_v", + "__et_offgraph_kv_layer_1_k", + "__et_offgraph_kv_layer_1_v", + "__et_offgraph_kv_cells", + "__et_offgraph_kv_mask_w0"}, + {}, + 0}; + declare_pool(container); + return container; + } +}; + +} // namespace + +TEST_F(CudaKVPoolTest, SideBuffersBindAtFixedAddressesAcrossGrowth) { + auto pool = make_pool(4); + auto container = full_container(); + auto handle = make_handle(container); + ASSERT_TRUE(pool->note_handle(&handle).get()); + ASSERT_EQ(pool->validate(), Error::Ok); + + ASSERT_EQ(pool->prepare(3, 0, cudaStreamPerThread), Error::Ok); + ASSERT_EQ(pool->bind(&handle), Error::Ok); + void* grow_k = container.bound["grow_k"].data; + void* fixed_k = container.bound["fixed_k"].data; + void* cells = container.bound["cells"].data; + void* mask = container.bound["mask"].data; + EXPECT_EQ(cells, pool->side_buffer(0)); + EXPECT_EQ(mask, pool->side_buffer(1)); + EXPECT_EQ(pool->side_buffer_bytes(0), 16 * sizeof(int64_t)); + EXPECT_EQ(container.bound["cells"].sizes, std::vector({16})); + const std::vector mask_sizes{1, 1, 16, kDeclaredRows}; + EXPECT_EQ(container.bound["mask"].sizes, mask_sizes); + EXPECT_EQ(container.bound["mask"].dtype, slimc10::ScalarType::Bool); + const std::vector fixed_sizes{1, kFixedRows, kHeads, kDim}; + EXPECT_EQ(container.bound["fixed_k"].sizes, fixed_sizes); + + // Side buffers start zeroed. + std::vector cells_host(16, -1); + ASSERT_EQ( + cudaMemcpy( + cells_host.data(), + cells, + cells_host.size() * sizeof(int64_t), + cudaMemcpyDeviceToHost), + cudaSuccess); + EXPECT_EQ(cells_host, std::vector(16, 0)); + ASSERT_EQ(pool->mark_step_done(), Error::Ok); + + // Growth moves only the growable layer. + ASSERT_EQ(pool->prepare(5, 3, cudaStreamPerThread), Error::Ok); + EXPECT_EQ(pool->rows(), 8); + EXPECT_EQ(pool->growth_count(), 1); + ASSERT_EQ(pool->bind(&handle), Error::Ok); + EXPECT_NE(container.bound["grow_k"].data, grow_k); + EXPECT_EQ(container.bound["fixed_k"].data, fixed_k); + EXPECT_EQ(container.bound["cells"].data, cells); + EXPECT_EQ(container.bound["mask"].data, mask); + pool->forget_handle(&handle); +} + +TEST_F(CudaKVPoolTest, GrowthCarriesLiveRowsAndCapsAtDeclaredRows) { + auto pool = make_pool(4); + auto container = full_container(); + auto handle = make_handle(container); + ASSERT_TRUE(pool->note_handle(&handle).get()); + ASSERT_EQ(pool->prepare(4, 0, cudaStreamPerThread), Error::Ok); + ASSERT_EQ(pool->bind(&handle), Error::Ok); + const auto history = iota_rows(3); + ASSERT_EQ( + cudaMemcpy( + container.bound["grow_k"].data, + history.data(), + history.size() * sizeof(uint16_t), + cudaMemcpyHostToDevice), + cudaSuccess); + ASSERT_EQ(pool->mark_step_done(), Error::Ok); + + // Asking past twice the rows grows straight to what is asked; past the + // declared rows it is capped there. + ASSERT_EQ(pool->prepare(20, 3, cudaStreamPerThread), Error::Ok); + EXPECT_EQ(pool->rows(), 20); + ASSERT_EQ(pool->bind(&handle), Error::Ok); + EXPECT_EQ(read_rows(container.bound["grow_k"].data, 3), history); + ASSERT_EQ(pool->prepare(kDeclaredRows, 3, cudaStreamPerThread), Error::Ok); + EXPECT_EQ(pool->rows(), kDeclaredRows); + EXPECT_EQ( + pool->allocated_bytes(), + 2 * (kDeclaredRows + kFixedRows) * kRow * + static_cast(sizeof(uint16_t))); + pool->forget_handle(&handle); +} + +TEST_F(CudaKVPoolTest, CompiledSizeRoundedUpBy64IsAccepted) { + // AOTI reports constant bytes rounded up to 64 when the program also holds + // CPU constants: the 128-byte cells buffer stays 128, but a 16-byte one + // reads as 64. Both forms describe the same declared shape. + cu::CudaKVPool pool( + {{kHeads, kDim, kDeclaredRows, /*growable=*/true}}, + {{"__et_offgraph_kv_read_len", slimc10::ScalarType::Long, {2}}}, + slimc10::ScalarType::BFloat16, + 4); + FakeContainer container{ + {"k", "v", "read_len"}, + {"__et_offgraph_kv_layer_0_k", + "__et_offgraph_kv_layer_0_v", + "__et_offgraph_kv_read_len"}, + {}, + 0}; + for (const char* fqn : + {"__et_offgraph_kv_layer_0_k", "__et_offgraph_kv_layer_0_v"}) { + declare( + container, + fqn, + slimc10::ScalarType::BFloat16, + {1, kDeclaredRows, kHeads, kDim}); + } + container.declared["__et_offgraph_kv_read_len"] = { + static_cast(slimc10::ScalarType::Long), 64}; + auto handle = make_handle(container); + EXPECT_TRUE(pool.note_handle(&handle).get()); + + // Rounded past the next multiple of 64 is a different shape. + container.declared["__et_offgraph_kv_read_len"].second = 128; + EXPECT_EQ(pool.note_handle(&handle).error(), Error::InvalidProgram); + pool.forget_handle(&handle); +} + +TEST_F(CudaKVPoolTest, ConstantNotMatchingItsCompiledSizeIsRejected) { + // A side buffer compiled at another shape than the pool allocates would + // be addressed past its allocation. + auto pool = make_pool(4); + auto container = full_container(); + declare( + container, + "__et_offgraph_kv_mask_w0", + slimc10::ScalarType::Bool, + {1, 1, 16, 2 * kDeclaredRows}); + auto handle = make_handle(container); + EXPECT_EQ(pool->note_handle(&handle).error(), Error::InvalidProgram); + + // And one compiled with another dtype. + auto wrong_dtype = full_container(); + declare(wrong_dtype, "__et_offgraph_kv_cells", slimc10::ScalarType::Int, {32}); + auto wrong_dtype_handle = make_handle(wrong_dtype); + EXPECT_EQ( + pool->note_handle(&wrong_dtype_handle).error(), Error::InvalidProgram); +} + +TEST_F(CudaKVPoolTest, ProgramMissingSideBuffersIsRejected) { + auto pool = make_pool(4); + // Every layer, but only one of the two side buffers: the lowering and the + // runtime disagree about the layout. + FakeContainer partial{ + {"grow_k", "grow_v", "fixed_k", "fixed_v", "cells"}, + {"__et_offgraph_kv_layer_0_k", + "__et_offgraph_kv_layer_0_v", + "__et_offgraph_kv_layer_1_k", + "__et_offgraph_kv_layer_1_v", + "__et_offgraph_kv_cells"}, + {}, + 0}; + declare_pool(partial); + auto handle = make_handle(partial); + EXPECT_EQ(pool->note_handle(&handle).error(), Error::InvalidProgram); + + // A program with no storage at all is simply not served. + FakeContainer embedding{{"weight"}, {"tok_embeddings.weight"}, {}, 0}; + auto embedding_handle = make_handle(embedding); + const auto serves = pool->note_handle(&embedding_handle); + ASSERT_EQ(serves.error(), Error::Ok); + EXPECT_FALSE(serves.get()); + // Neither program contributed storage, the rejected one included. + EXPECT_EQ(pool->validate(), Error::InvalidProgram); +} + +TEST_F(CudaKVPoolTest, FailedSideBufferAllocationLeavesThePoolRetryable) { + // A side buffer too large for any device fails the first allocation, which + // must release the layers it already allocated. + cu::CudaKVPool pool( + {{kHeads, kDim, kDeclaredRows, /*growable=*/true}}, + {{"__et_offgraph_kv_huge", + slimc10::ScalarType::Byte, + {int64_t{1} << 50}}}, + slimc10::ScalarType::BFloat16, + 4); + FakeContainer container{ + {"k", "v", "huge"}, + {"__et_offgraph_kv_layer_0_k", + "__et_offgraph_kv_layer_0_v", + "__et_offgraph_kv_huge"}, + {}, + 0}; + for (const char* fqn : + {"__et_offgraph_kv_layer_0_k", "__et_offgraph_kv_layer_0_v"}) { + declare( + container, + fqn, + slimc10::ScalarType::BFloat16, + {1, kDeclaredRows, kHeads, kDim}); + } + declare( + container, + "__et_offgraph_kv_huge", + slimc10::ScalarType::Byte, + {int64_t{1} << 50}); + auto handle = make_handle(container); + ASSERT_TRUE(pool.note_handle(&handle).get()); + EXPECT_NE(pool.prepare(1, 0, cudaStreamPerThread), Error::Ok); + EXPECT_FALSE(pool.allocated()); + EXPECT_EQ(pool.allocated_bytes(), 0); + pool.forget_handle(&handle); +}