diff --git a/ggml/src/ggml-vulkan/ggml-vulkan.cpp b/ggml/src/ggml-vulkan/ggml-vulkan.cpp index 19e7fbdaae..647de05aa0 100644 --- a/ggml/src/ggml-vulkan/ggml-vulkan.cpp +++ b/ggml/src/ggml-vulkan/ggml-vulkan.cpp @@ -580,6 +580,17 @@ static constexpr std::initializer_list> rms_norm_mul_rope_vie }; +struct vk_peer_staging { + void * host_ptr = nullptr; + size_t host_size = 0; + vk_buffer src_buf; + vk_buffer dst_buf; + // sync_fd semaphores for GPU-only cross-device synchronization + vk::Semaphore src_sem; // exportable binary semaphore on source device + vk::Semaphore dst_sem; // importable binary semaphore on dest device + bool use_sync_fd = false; +}; + struct vk_device_struct { std::recursive_mutex mutex; @@ -591,6 +602,7 @@ struct vk_device_struct { uint64_t suballocation_block_size; uint64_t min_imported_host_pointer_alignment; bool external_memory_host {}; + bool external_semaphore_sync_fd {}; bool fp16; bool bf16; bool pipeline_robustness; @@ -857,6 +869,8 @@ struct vk_device_struct { vk::Fence fence; vk_buffer sync_staging; + std::map peer_staging; + ggml_backend_buffer_type buffer_type; bool disable_fusion; @@ -871,6 +885,25 @@ struct vk_device_struct { device.destroyFence(fence); + for (auto& [peer, staging] : peer_staging) { + if (staging.src_sem) { + device.destroySemaphore(staging.src_sem); + } + if (staging.dst_sem) { + peer->device.destroySemaphore(staging.dst_sem); + } + staging.src_buf.reset(); + staging.dst_buf.reset(); + if (staging.host_ptr) { +#if defined(_MSC_VER) || defined(__MINGW32__) + _aligned_free(staging.host_ptr); +#else + free(staging.host_ptr); +#endif + } + } + peer_staging.clear(); + ggml_vk_destroy_buffer(sync_staging); compute_queue.cmd_pool.destroy(device); @@ -4882,6 +4915,8 @@ static vk_device ggml_vk_get_device(size_t idx) { device->memory_priority = true; } else if (strcmp("VK_EXT_external_memory_host", properties.extensionName) == 0) { device->external_memory_host = true; + } else if (strcmp("VK_KHR_external_semaphore_fd", properties.extensionName) == 0) { + device->external_semaphore_sync_fd = true; #if defined(VK_EXT_shader_64bit_indexing) } else if (strcmp("VK_EXT_shader_64bit_indexing", properties.extensionName) == 0) { device->shader_64b_indexing = true; @@ -5181,6 +5216,22 @@ static vk_device ggml_vk_get_device(size_t idx) { device_extensions.push_back("VK_EXT_external_memory_host"); } + if (device->external_semaphore_sync_fd) { + // Verify SYNC_FD_BIT is actually supported for binary semaphores + vk::PhysicalDeviceExternalSemaphoreInfo sem_info{ + vk::ExternalSemaphoreHandleTypeFlagBits::eSyncFd + }; + vk::ExternalSemaphoreProperties sem_props = + device->physical_device.getExternalSemaphoreProperties(sem_info); + if ((sem_props.externalSemaphoreFeatures & vk::ExternalSemaphoreFeatureFlagBits::eExportable) && + (sem_props.externalSemaphoreFeatures & vk::ExternalSemaphoreFeatureFlagBits::eImportable)) { + device_extensions.push_back("VK_KHR_external_semaphore_fd"); + device_extensions.push_back("VK_KHR_external_semaphore"); + } else { + device->external_semaphore_sync_fd = false; + } + } + #if defined(VK_EXT_shader_64bit_indexing) VkPhysicalDeviceShader64BitIndexingFeaturesEXT shader_64bit_indexing_features {}; shader_64bit_indexing_features.sType = VK_STRUCTURE_TYPE_PHYSICAL_DEVICE_SHADER_64_BIT_INDEXING_FEATURES_EXT; @@ -13757,6 +13808,96 @@ static void ggml_backend_vk_get_tensor_async(ggml_backend_t backend, const ggml_ } } +static vk_buffer ggml_vk_buffer_from_host_ptr(vk_device & device, void * ptr, size_t size); + +static bool ggml_vk_ensure_peer_staging(vk_device& src_dev, vk_device& dst_dev, size_t required_size) { + if (!src_dev->external_memory_host || !dst_dev->external_memory_host) { + return false; + } + + auto it = src_dev->peer_staging.find(dst_dev.get()); + if (it != src_dev->peer_staging.end() && it->second.host_size >= required_size) { + return true; + } + + // Tear down old entry if too small + if (it != src_dev->peer_staging.end()) { + if (it->second.src_sem) { + src_dev->device.destroySemaphore(it->second.src_sem); + } + if (it->second.dst_sem) { + dst_dev->device.destroySemaphore(it->second.dst_sem); + } + it->second.src_buf.reset(); + it->second.dst_buf.reset(); + if (it->second.host_ptr) { +#if defined(_MSC_VER) || defined(__MINGW32__) + _aligned_free(it->second.host_ptr); +#else + free(it->second.host_ptr); +#endif + } + src_dev->peer_staging.erase(it); + } + + uint64_t alignment = std::max(src_dev->min_imported_host_pointer_alignment, + dst_dev->min_imported_host_pointer_alignment); + if (alignment == 0) { + alignment = 4096; + } + + size_t alloc_size = CEIL_DIV(required_size, alignment) * alignment; + + void * host_ptr = nullptr; +#if defined(_MSC_VER) || defined(__MINGW32__) + host_ptr = _aligned_malloc(alloc_size, (size_t)alignment); +#else + if (posix_memalign(&host_ptr, (size_t)alignment, alloc_size) != 0) { + host_ptr = nullptr; + } +#endif + if (!host_ptr) { + return false; + } + + vk_buffer src_buf = ggml_vk_buffer_from_host_ptr(src_dev, host_ptr, alloc_size); + if (!src_buf) { +#if defined(_MSC_VER) || defined(__MINGW32__) + _aligned_free(host_ptr); +#else + free(host_ptr); +#endif + return false; + } + + vk_buffer dst_buf = ggml_vk_buffer_from_host_ptr(dst_dev, host_ptr, alloc_size); + if (!dst_buf) { + src_buf.reset(); +#if defined(_MSC_VER) || defined(__MINGW32__) + _aligned_free(host_ptr); +#else + free(host_ptr); +#endif + return false; + } + + vk::Semaphore src_sem{}, dst_sem{}; + bool use_sync_fd = src_dev->external_semaphore_sync_fd && dst_dev->external_semaphore_sync_fd; + if (use_sync_fd) { + vk::ExportSemaphoreCreateInfo export_ci{ vk::ExternalSemaphoreHandleTypeFlagBits::eSyncFd }; + vk::SemaphoreCreateInfo sci{}; + sci.setPNext(&export_ci); + src_sem = src_dev->device.createSemaphore(sci); + + vk::SemaphoreCreateInfo dst_sci{}; + dst_sem = dst_dev->device.createSemaphore(dst_sci); + } + + src_dev->peer_staging[dst_dev.get()] = { host_ptr, alloc_size, src_buf, dst_buf, + src_sem, dst_sem, use_sync_fd }; + return true; +} + static bool ggml_backend_vk_cpy_tensor_async(ggml_backend_t backend_src, ggml_backend_t backend_dst, const ggml_tensor * src, ggml_tensor * dst) { VK_LOG_DEBUG("ggml_backend_vk_cpy_tensor_async(" << src << " -> " << dst << ", size=" << ggml_nbytes(src) << ")"); ggml_backend_vk_context * ctx = (ggml_backend_vk_context *)backend_dst->context; @@ -13776,9 +13917,87 @@ static bool ggml_backend_vk_cpy_tensor_async(ggml_backend_t backend_src, ggml_ba if (ggml_backend_buffer_is_vk(src->buffer)) { ggml_backend_vk_buffer_context * src_buf_ctx = (ggml_backend_vk_buffer_context *)src->buffer->context; - // Async copy only works within the same device if (src_buf_ctx->dev_buffer->device != dst_buf->device) { - return false; + // Cross-device copy via shared staging buffer + vk_device src_dev = src_buf_ctx->dev_buffer->device; + vk_device dst_dev = ctx->device; + + size_t nbytes = ggml_nbytes(src); + vk_buffer src_vk_buf = src_buf_ctx->dev_buffer; + size_t src_offset = vk_tensor_offset(src) + src->view_offs; + size_t dst_offset = vk_tensor_offset(dst) + dst->view_offs; + + if (!ggml_vk_ensure_peer_staging(src_dev, dst_dev, nbytes)) { + return false; + } + + auto& staging = src_dev->peer_staging[dst_dev.get()]; + + // HOP 1: src VRAM → shared staging (on source compute queue) + // Submitted on the compute queue — implicit queue submission ordering + // guarantees this executes after all prior compute work. + vk_context hop1_ctx; + { + std::lock_guard guard(src_dev->mutex); + hop1_ctx = ggml_vk_create_temporary_context(src_dev->compute_queue.cmd_pool); + ggml_vk_ctx_begin(src_dev, hop1_ctx); + + VkBufferCopy bc{ src_offset, 0, nbytes }; + vkCmdCopyBuffer(hop1_ctx->s->buffer->buf, + (VkBuffer)src_vk_buf->buffer, + (VkBuffer)staging.src_buf->buffer, + 1, &bc); + + ggml_vk_ctx_end(hop1_ctx); + } + + // HOP 2: shared staging → dst VRAM (on dest device) + vk_context dst_compute_ctx = ggml_vk_get_compute_ctx(ctx); + + if (staging.use_sync_fd) { + // Tier 1: GPU-only synchronization via sync_fd + // Signal exportable semaphore after hop1, submit without fence + hop1_ctx->seqs.back().back().signal_semaphores.push_back({ staging.src_sem, 0 }); + ggml_vk_submit(hop1_ctx, {}); + + // Export sync_fd from source semaphore + vk::SemaphoreGetFdInfoKHR get_fd_info{ + staging.src_sem, + vk::ExternalSemaphoreHandleTypeFlagBits::eSyncFd + }; + int sync_fd = src_dev->device.getSemaphoreFdKHR(get_fd_info); + + // Import sync_fd into destination semaphore + vk::ImportSemaphoreFdInfoKHR import_info{ + staging.dst_sem, + vk::SemaphoreImportFlagBits::eTemporary, + vk::ExternalSemaphoreHandleTypeFlagBits::eSyncFd, + sync_fd + }; + dst_dev->device.importSemaphoreFdKHR(import_info); + + // Destination waits on imported semaphore before hop2 + dst_compute_ctx->s->wait_semaphores.push_back({ staging.dst_sem, 0 }); + } else { + // Tier 2: CPU fence fallback + // Submit hop1 with fence, wait for just this transfer + { + std::lock_guard guard(src_dev->mutex); + ggml_vk_submit(hop1_ctx, src_dev->fence); + VK_CHECK(src_dev->device.waitForFences({ src_dev->fence }, true, UINT64_MAX), + "cross_device_hop1 waitForFences"); + src_dev->device.resetFences({ src_dev->fence }); + ggml_vk_queue_command_pools_cleanup(src_dev); + } + } + + VkBufferCopy bc2{ 0, dst_offset, nbytes }; + vkCmdCopyBuffer(dst_compute_ctx->s->buffer->buf, + (VkBuffer)staging.dst_buf->buffer, + (VkBuffer)dst_buf->buffer, + 1, &bc2); + + return true; } vk_context compute_ctx = ggml_vk_get_compute_ctx(ctx);