https://github.com/jellytabby updated https://github.com/llvm/llvm-project/pull/218049
>From e488e1d0bcde734ecd71ce8b937a642ca0488c08 Mon Sep 17 00:00:00 2001 From: Sophia Herrmann <[email protected]> Date: Thu, 13 Aug 2026 16:47:23 -0700 Subject: [PATCH 1/2] add blocking semantics to LaunchKernel and Memcpy --- clang/lib/CodeGen/CGCUDANV.cpp | 4 +- clang/lib/Driver/ToolChains/Clang.cpp | 1 + .../languages/kernel/include/LanguageUtils.h | 50 ++++++- offload/languages/kernel/include/State.h | 2 +- .../languages/kernel/src/LanguageLaunch.cpp | 20 ++- .../languages/kernel/src/LanguageRuntime.cpp | 15 +- .../CUDA/blocking_stream_semantics.cu | 131 ++++++++++++++++++ offload/test/offloading/CUDA/stream_api.cu | 6 +- .../HIP/blocking_stream_semantics.hip | 125 +++++++++++++++++ offload/test/offloading/HIP/stream_api.hip | 6 +- 10 files changed, 342 insertions(+), 18 deletions(-) create mode 100644 offload/test/offloading/CUDA/blocking_stream_semantics.cu create mode 100644 offload/test/offloading/HIP/blocking_stream_semantics.hip diff --git a/clang/lib/CodeGen/CGCUDANV.cpp b/clang/lib/CodeGen/CGCUDANV.cpp index e03b7e754ab3f..14904c52b98ff 100644 --- a/clang/lib/CodeGen/CGCUDANV.cpp +++ b/clang/lib/CodeGen/CGCUDANV.cpp @@ -442,7 +442,9 @@ void CGNVCUDARuntime::emitDeviceStubBodyNew(CodeGenFunction &CGF, std::string KernelLaunchAPI = "LaunchKernel"; if (CGF.getLangOpts().GPUDefaultStream == LangOptions::GPUDefaultStreamKind::PerThread) { - if (CGF.getLangOpts().HIP) + if (CGF.getLangOpts().OffloadViaLLVM) + KernelLaunchAPI = KernelLaunchAPI + ""; + else if (CGF.getLangOpts().HIP) KernelLaunchAPI = KernelLaunchAPI + "_spt"; else if (CGF.getLangOpts().CUDA) KernelLaunchAPI = KernelLaunchAPI + "_ptsz"; diff --git a/clang/lib/Driver/ToolChains/Clang.cpp b/clang/lib/Driver/ToolChains/Clang.cpp index 54583fe3abbd8..8b21400ab959f 100644 --- a/clang/lib/Driver/ToolChains/Clang.cpp +++ b/clang/lib/Driver/ToolChains/Clang.cpp @@ -8396,6 +8396,7 @@ void Clang::ConstructJob(Compilation &C, const JobAction &JA, if (IsHIP) { CmdArgs.push_back("-fcuda-allow-variadic-functions"); + /// TODO: Why is this not forwarded when IsCUDA? Args.AddLastArg(CmdArgs, options::OPT_fgpu_default_stream_EQ); } diff --git a/offload/languages/kernel/include/LanguageUtils.h b/offload/languages/kernel/include/LanguageUtils.h index f0ced30a1287d..708de92fcdab6 100644 --- a/offload/languages/kernel/include/LanguageUtils.h +++ b/offload/languages/kernel/include/LanguageUtils.h @@ -14,6 +14,10 @@ #include "State.h" #include "Stream.h" +using RuntimeState = llvm::offload::StateTy; +using ThreadState = llvm::offload::ThreadStateTy; +using StreamTy = llvm::offload::StreamTy; + namespace llvm { namespace offload { @@ -43,8 +47,7 @@ static inline Error_t convertResult(ol_result_t Result) { /// Set the last error for the current thread and return it. static inline Error_t setLastError(Error_t Error) { // TODO: find a more efficient way to set last error - return static_cast<Error_t>( - llvm::offload::ThreadStateTy::setLastError(Error)); + return static_cast<Error_t>(ThreadState::setLastError(Error)); } /// Convert an ol_result_t to the active language's Error_t and set it as the @@ -54,12 +57,51 @@ static inline Error_t convertAndSetLastError(ol_result_t Result) { } /// Convert between the language-facing opaque stream and the internal stream. -static inline Stream_t toLanguageStream(llvm::offload::StreamTy *Stream) { +static inline Stream_t toLanguageStream(StreamTy *Stream) { return reinterpret_cast<Stream_t>(Stream); } static inline llvm::offload::StreamTy *toInternalStream(Stream_t Stream) { - return reinterpret_cast<llvm::offload::StreamTy *>(Stream); + return reinterpret_cast<StreamTy *>(Stream); +} + +/// Wait for blocking streams before executing if we are legacy default stream. +static inline ol_result_t waitOnBlockingStreams() { + ol_device_handle_t Device = ThreadState::getDefaultDevice(); + if (!RuntimeState::hasLegacyDefaultStream(Device) || + RuntimeState::getBlockingStreams(Device).empty()) + return OL_SUCCESS; + StreamTy *DefaultStream = ThreadState::getDefaultStream(); + llvm::SmallVector<ol_event_handle_t, 8> Events; + for (StreamTy *BlockingStream : RuntimeState::getBlockingStreams(Device)) { + ol_event_handle_t Event = nullptr; + ol_result_t Result = + olCreateEvent(BlockingStream->Queue, OL_EVENT_FLAGS_NONE, &Event); + if (Result != OL_SUCCESS) + return Result; + Events.push_back(Event); + } + + return olWaitEvents(DefaultStream->Queue, Events.data(), Events.size()); +} + +/// Wait for the legacy default stream to complete before launching a kernel on +/// a blocking stream. +static inline ol_result_t waitOnLegacyDefaultStream(StreamTy *SourceStream, + ol_device_handle_t Device) { + if (!RuntimeState::hasLegacyDefaultStream(Device)) + return OL_SUCCESS; + + StreamTy *DefaultStream = ThreadState::getDefaultStream(); + assert(DefaultStream->Kind == llvm::offload::QueueKind::LegacyDefault && + "Default stream is not a legacy default stream"); + + ol_event_handle_t Event = nullptr; + ol_result_t Result = + olCreateEvent(DefaultStream->Queue, OL_EVENT_FLAGS_NONE, &Event); + if (Result != OL_SUCCESS) + return Result; + return olWaitEvents(SourceStream->Queue, &Event, 1); } /// Convert a Stream_t to an ol_queue_handle_t. diff --git a/offload/languages/kernel/include/State.h b/offload/languages/kernel/include/State.h index bb5422b5d3bb6..e7c45bc8dec88 100644 --- a/offload/languages/kernel/include/State.h +++ b/offload/languages/kernel/include/State.h @@ -148,7 +148,7 @@ struct StateTy { static llvm::SmallPtrSet<StreamTy *, 8> getBlockingStreams(ol_device_handle_t Device); - /// Return true if \p Device has an existing legacy default stream. + /// Return true if \p Device has an initialized (i.e. previously used) legacy default stream. static bool hasLegacyDefaultStream(ol_device_handle_t Device); /// Create a stream for \p Device and register it with the process state. diff --git a/offload/languages/kernel/src/LanguageLaunch.cpp b/offload/languages/kernel/src/LanguageLaunch.cpp index 48505ed376e43..a305e42cac0de 100644 --- a/offload/languages/kernel/src/LanguageLaunch.cpp +++ b/offload/languages/kernel/src/LanguageLaunch.cpp @@ -8,6 +8,7 @@ #include "LanguageLaunch.h" #include "LanguageUtils.h" +#include "OffloadAPI.h" #include "OffloadErrors.h" #include "State.h" #include "Stream.h" @@ -51,8 +52,21 @@ ol_result_t __llvmLaunchKernelImpl(const char *KernelID, dim3 GridDim, LaunchSizeArgs.GroupSize.z = BlockDim.z; LaunchSizeArgs.DynSharedMemory = DynamicSharedMem; - ol_queue_handle_t Queue = Stream ? reinterpret_cast<StreamTy *>(Stream)->Queue - : ThreadState::getDefaultQueue(); + StreamTy *LaunchStream = Stream ? reinterpret_cast<StreamTy *>(Stream) + : ThreadState::getDefaultStream(); + if (!LaunchStream || !RuntimeState::isStreamRegistered(LaunchStream) || + LaunchStream->Device != Device) + return &InvalidConfigurationError; + + if (LaunchStream->Kind == llvm::offload::QueueKind::LegacyDefault) { + ol_result_t Result = waitOnBlockingStreams(); + if (Result != OL_SUCCESS) + return Result; + } else if (LaunchStream->Kind == llvm::offload::QueueKind::ExplicitBlocking) { + ol_result_t Result = waitOnLegacyDefaultStream(LaunchStream, Device); + if (Result != OL_SUCCESS) + return Result; + } struct OffloadKernelArgs { void **Args; @@ -68,7 +82,7 @@ ol_result_t __llvmLaunchKernelImpl(const char *KernelID, dim3 GridDim, if (!OKA->Args[I] || OKA->ArgSizes[I] == 0) return &InvalidArgumentError; - return olLaunchKernel(Queue, Device, Kernel, &LaunchSizeArgs, + return olLaunchKernel(LaunchStream->Queue, Device, Kernel, &LaunchSizeArgs, /*Properties=*/nullptr, OKA->NumArgs, OKA->Args, OKA->ArgSizes); } diff --git a/offload/languages/kernel/src/LanguageRuntime.cpp b/offload/languages/kernel/src/LanguageRuntime.cpp index c387ada6661c2..a57d10688b644 100644 --- a/offload/languages/kernel/src/LanguageRuntime.cpp +++ b/offload/languages/kernel/src/LanguageRuntime.cpp @@ -23,6 +23,7 @@ #include "Types.h" #include "OffloadAPI.h" +#include "llvm/ADT/SmallVector.h" #include <cassert> #include <cstdio> @@ -51,8 +52,13 @@ Error_t Free(void *DevPtr) { } Error_t Memcpy(void *Dst, const void *Src, size_t Size, MemcpyKind Kind) { + if (Kind != MemcpyHostToHost) { + ol_result_t Result = waitOnBlockingStreams(); + if (Result != OL_SUCCESS) + return convertAndSetLastError(Result); + } + ol_device_handle_t Device = ThreadState::getDefaultDevice(); ol_queue_handle_t Queue = ThreadState::getDefaultQueue(); - ol_result_t Result; switch (Kind) { case MemcpyHostToHost: { @@ -61,21 +67,17 @@ Error_t Memcpy(void *Dst, const void *Src, size_t Size, MemcpyKind Kind) { break; } case MemcpyHostToDevice: { - ol_device_handle_t Device = ThreadState::getDefaultDevice(); ol_device_handle_t Host = RuntimeState::getHostDevice(); Result = olMemcpy(Queue, Dst, Device, const_cast<void *>(Src), Host, Size); break; } case MemcpyDeviceToHost: { - ol_device_handle_t Device = ThreadState::getDefaultDevice(); ol_device_handle_t Host = RuntimeState::getHostDevice(); Result = olMemcpy(Queue, Dst, Host, const_cast<void *>(Src), Device, Size); break; } case MemcpyDeviceToDevice: { - ol_device_handle_t Device = ThreadState::getDefaultDevice(); - Result = olMemcpy(Queue, Dst, Device, const_cast<void *>(Src), Device, Size); break; @@ -87,6 +89,9 @@ Error_t Memcpy(void *Dst, const void *Src, size_t Size, MemcpyKind Kind) { if (Result != OL_SUCCESS) return convertAndSetLastError(Result); + if (!Queue) + return convertAndSetLastError(Result); + Result = olSyncQueue(Queue); return convertAndSetLastError(Result); } diff --git a/offload/test/offloading/CUDA/blocking_stream_semantics.cu b/offload/test/offloading/CUDA/blocking_stream_semantics.cu new file mode 100644 index 0000000000000..83aa3383571f6 --- /dev/null +++ b/offload/test/offloading/CUDA/blocking_stream_semantics.cu @@ -0,0 +1,131 @@ +// clang-format off +// RUN: %clang++ %flags -foffload-via-llvm --offload-arch=native %s -o %t +// RUN: %t | %fcheck-generic --check-prefix=LEGACY +// RUN: %clang++ %flags -foffload-via-llvm --offload-arch=native %s -o %t -fopenmp +// RUN: %t | %fcheck-generic --check-prefix=LEGACY +// RUN: %clang++ %flags -foffload-via-llvm --offload-arch=native %s -o %t -fgpu-default-stream=per-thread +// RUN: %t | %fcheck-generic --check-prefix=PERTHREAD +// clang-format on + +// UNSUPPORTED: aarch64-unknown-linux-gnu +// UNSUPPORTED: x86_64-unknown-linux-gnu +// UNSUPPORTED: nvptx64-nvidia-cuda-LTO +// UNSUPPORTED: amdgcn-amd-amdhsa-LTO +// UNSUPPORTED: amdgpu-amd-amdhsa-LTO +// UNSUPPORTED: intelgpu + +#include <stdio.h> + +__global__ void delayedSetValue(int *Out, int Value) { + volatile unsigned long long Delay = 0; + for (unsigned I = 0; I < 1000000; ++I) + Delay += I; + if (Delay) + *Out = Value; +} + +__global__ void copyValue(int *In, int *Out) { *Out = *In; } + +__global__ void waitThenSetValue(int *Gate, int *Out, int Value) { + volatile int *VolatileGate = Gate; + for (unsigned I = 0; I < 100000000 && *VolatileGate == 0; ++I) + ; + *Out = Value; +} + +__global__ void copyValueAndRelease(int *In, int *Out, int *Gate) { + *Out = *In; + volatile int *VolatileGate = Gate; + *VolatileGate = 1; +} + +int main(int argc, char **argv) { + cudaStream_t BlockingStream = nullptr; + if (cudaStreamCreateWithFlags(&BlockingStream, cudaStreamDefault) != + cudaSuccess) + return 1; + cudaStream_t NonBlockingStream = nullptr; + if (cudaStreamCreateWithFlags(&NonBlockingStream, cudaStreamNonBlocking) != + cudaSuccess) + return 1; + + int *In = nullptr; + int *Out = nullptr; + int *Gate = nullptr; + if (cudaMalloc(&In, sizeof(int)) != cudaSuccess) + return 1; + if (cudaMalloc(&Out, sizeof(int)) != cudaSuccess) + return 1; + if (cudaMalloc(&Gate, sizeof(int)) != cudaSuccess) + return 1; + + int Initial = 0; + int Result = 0; + if (cudaMemcpy(In, &Initial, sizeof(int), cudaMemcpyHostToDevice) != + cudaSuccess) + return 1; + if (cudaMemcpy(Out, &Initial, sizeof(int), cudaMemcpyHostToDevice) != + cudaSuccess) + return 1; + + delayedSetValue<<<1, 1, 0, BlockingStream>>>(In, 99); + copyValue<<<1, 1>>>(In, Out); + if (cudaMemcpy(&Result, Out, sizeof(int), cudaMemcpyDeviceToHost) != + cudaSuccess) + return 1; + + printf("legacy default waited on blocking stream: %d\n", Result); + // LEGACY: legacy default waited on blocking stream: 99 + // PERTHREAD: legacy default waited on blocking stream: 0 + + Result = 0; + if (cudaMemcpy(Out, &Initial, sizeof(int), cudaMemcpyHostToDevice) != + cudaSuccess) + return 1; + + delayedSetValue<<<1, 1>>>(In, 123); + copyValue<<<1, 1, 0, BlockingStream>>>(In, Out); + if (cudaStreamSynchronize(BlockingStream) != cudaSuccess) + return 1; + if (cudaMemcpy(&Result, Out, sizeof(int), cudaMemcpyDeviceToHost) != + cudaSuccess) + return 1; + + printf("blocking stream waited on legacy default: %d\n", Result); + // LEGACY: blocking stream waited on legacy default: 123 + // PERTHREAD: blocking stream waited on legacy default: 99 + + Result = 0; + if (cudaMemcpy(In, &Initial, sizeof(int), cudaMemcpyHostToDevice) != + cudaSuccess) + return 1; + if (cudaMemcpy(Out, &Initial, sizeof(int), cudaMemcpyHostToDevice) != + cudaSuccess) + return 1; + if (cudaMemcpy(Gate, &Initial, sizeof(int), cudaMemcpyHostToDevice) != + cudaSuccess) + return 1; + + waitThenSetValue<<<1, 1>>>(Gate, In, 321); + copyValueAndRelease<<<1, 1, 0, NonBlockingStream>>>(In, Out, Gate); + if (cudaStreamSynchronize(NonBlockingStream) != cudaSuccess) + return 1; + if (cudaMemcpy(&Result, Out, sizeof(int), cudaMemcpyDeviceToHost) != + cudaSuccess) + return 1; + + printf("nonblocking stream did not wait on legacy default: %d\n", Result); + // LEGACY: nonblocking stream did not wait on legacy default: 0 + // PERTHREAD: nonblocking stream did not wait on legacy default: 0 + + if (cudaStreamDestroy(BlockingStream) != cudaSuccess) + return 1; + if (cudaStreamDestroy(NonBlockingStream) != cudaSuccess) + return 1; + if (cudaFree(In) != cudaSuccess) + return 1; + if (cudaFree(Out) != cudaSuccess) + return 1; + if (cudaFree(Gate) != cudaSuccess) + return 1; +} diff --git a/offload/test/offloading/CUDA/stream_api.cu b/offload/test/offloading/CUDA/stream_api.cu index c5b1caaa5315c..0e1c328cc5f1b 100644 --- a/offload/test/offloading/CUDA/stream_api.cu +++ b/offload/test/offloading/CUDA/stream_api.cu @@ -70,8 +70,6 @@ int main(int argc, char **argv) { if (cudaStreamSynchronize(Stream) != cudaSuccess) return 1; - if (cudaDeviceSynchronize() != cudaSuccess) - return 1; if (cudaMemcpy(&StreamResult, StreamPtr, sizeof(int), cudaMemcpyDeviceToHost) != cudaSuccess) return 1; @@ -86,6 +84,10 @@ int main(int argc, char **argv) { if (cudaStreamDestroy(Stream) != cudaSuccess) return 1; + if (cudaStreamDestroy(BlockingStream) != cudaSuccess) + return 1; + if (cudaStreamDestroy(NonBlockingStream) != cudaSuccess) + return 1; print_error("destroyed stream destroy", cudaStreamDestroy(Stream)); // CHECK: destroyed stream destroy value: 4 // CHECK: destroyed stream destroy name: cudaErrorInvalidResourceHandle diff --git a/offload/test/offloading/HIP/blocking_stream_semantics.hip b/offload/test/offloading/HIP/blocking_stream_semantics.hip new file mode 100644 index 0000000000000..8c28c7b1b2ee8 --- /dev/null +++ b/offload/test/offloading/HIP/blocking_stream_semantics.hip @@ -0,0 +1,125 @@ +// clang-format off +// RUN: %clang++ %flags -foffload-via-llvm --offload-arch=native %s -o %t +// RUN: %t | %fcheck-generic --check-prefix=LEGACY +// RUN: %clang++ %flags -foffload-via-llvm --offload-arch=native %s -o %t -fopenmp +// RUN: %t | %fcheck-generic --check-prefix=LEGACY +// RUN: %clang++ %flags -foffload-via-llvm --offload-arch=native %s -o %t -fgpu-default-stream=per-thread +// RUN: %t | %fcheck-generic --check-prefix=PERTHREAD +// clang-format on + +// UNSUPPORTED: aarch64-unknown-linux-gnu +// UNSUPPORTED: x86_64-unknown-linux-gnu +// UNSUPPORTED: nvptx64-nvidia-cuda-LTO +// UNSUPPORTED: amdgcn-amd-amdhsa-LTO +// UNSUPPORTED: amdgpu-amd-amdhsa-LTO +// UNSUPPORTED: intelgpu + +#include <stdio.h> + +__global__ void delayedSetValue(int *Out, int Value) { + volatile unsigned long long Delay = 0; + for (unsigned I = 0; I < 1000000; ++I) + Delay += I; + if (Delay) + *Out = Value; +} + +__global__ void copyValue(int *In, int *Out) { *Out = *In; } + +__global__ void waitThenSetValue(int *Gate, int *Out, int Value) { + volatile int *VolatileGate = Gate; + for (unsigned I = 0; I < 100000000 && *VolatileGate == 0; ++I) + ; + *Out = Value; +} + +__global__ void copyValueAndRelease(int *In, int *Out, int *Gate) { + *Out = *In; + volatile int *VolatileGate = Gate; + *VolatileGate = 1; +} + +int main(int argc, char **argv) { + hipStream_t BlockingStream = nullptr; + if (hipStreamCreateWithFlags(&BlockingStream, hipStreamDefault) != hipSuccess) + return 1; + hipStream_t NonBlockingStream = nullptr; + if (hipStreamCreateWithFlags(&NonBlockingStream, hipStreamNonBlocking) != + hipSuccess) + return 1; + + int *In = nullptr; + int *Out = nullptr; + int *Gate = nullptr; + if (hipMalloc(&In, sizeof(int)) != hipSuccess) + return 1; + if (hipMalloc(&Out, sizeof(int)) != hipSuccess) + return 1; + if (hipMalloc(&Gate, sizeof(int)) != hipSuccess) + return 1; + + int Initial = 0; + int Result = 0; + if (hipMemcpy(In, &Initial, sizeof(int), hipMemcpyHostToDevice) != hipSuccess) + return 1; + if (hipMemcpy(Out, &Initial, sizeof(int), hipMemcpyHostToDevice) != + hipSuccess) + return 1; + + delayedSetValue<<<1, 1, 0, BlockingStream>>>(In, 99); + copyValue<<<1, 1>>>(In, Out); + if (hipMemcpy(&Result, Out, sizeof(int), hipMemcpyDeviceToHost) != hipSuccess) + return 1; + + printf("legacy default waited on blocking stream: %d\n", Result); + // LEGACY: legacy default waited on blocking stream: 99 + // PERTHREAD: legacy default waited on blocking stream: 0 + + Result = 0; + if (hipMemcpy(Out, &Initial, sizeof(int), hipMemcpyHostToDevice) != + hipSuccess) + return 1; + + delayedSetValue<<<1, 1>>>(In, 123); + copyValue<<<1, 1, 0, BlockingStream>>>(In, Out); + if (hipStreamSynchronize(BlockingStream) != hipSuccess) + return 1; + if (hipMemcpy(&Result, Out, sizeof(int), hipMemcpyDeviceToHost) != hipSuccess) + return 1; + + printf("blocking stream waited on legacy default: %d\n", Result); + // LEGACY: blocking stream waited on legacy default: 123 + // PERTHREAD: blocking stream waited on legacy default: 99 + + Result = 0; + if (hipMemcpy(In, &Initial, sizeof(int), hipMemcpyHostToDevice) != hipSuccess) + return 1; + if (hipMemcpy(Out, &Initial, sizeof(int), hipMemcpyHostToDevice) != + hipSuccess) + return 1; + if (hipMemcpy(Gate, &Initial, sizeof(int), hipMemcpyHostToDevice) != + hipSuccess) + return 1; + + waitThenSetValue<<<1, 1>>>(Gate, In, 321); + copyValueAndRelease<<<1, 1, 0, NonBlockingStream>>>(In, Out, Gate); + if (hipStreamSynchronize(NonBlockingStream) != hipSuccess) + return 1; + if (hipMemcpy(&Result, Out, sizeof(int), hipMemcpyDeviceToHost) != hipSuccess) + return 1; + + printf("nonblocking stream did not wait on legacy default: %d\n", Result); + // LEGACY: nonblocking stream did not wait on legacy default: 0 + // PERTHREAD: nonblocking stream did not wait on legacy default: 0 + + if (hipStreamDestroy(BlockingStream) != hipSuccess) + return 1; + if (hipStreamDestroy(NonBlockingStream) != hipSuccess) + return 1; + if (hipFree(In) != hipSuccess) + return 1; + if (hipFree(Out) != hipSuccess) + return 1; + if (hipFree(Gate) != hipSuccess) + return 1; +} diff --git a/offload/test/offloading/HIP/stream_api.hip b/offload/test/offloading/HIP/stream_api.hip index fbfca230ee697..3460a6a0c351c 100644 --- a/offload/test/offloading/HIP/stream_api.hip +++ b/offload/test/offloading/HIP/stream_api.hip @@ -69,8 +69,6 @@ int main(int argc, char **argv) { if (hipStreamSynchronize(Stream) != hipSuccess) return 1; - if (hipDeviceSynchronize() != hipSuccess) - return 1; if (hipMemcpy(&StreamResult, StreamPtr, sizeof(int), hipMemcpyDeviceToHost) != hipSuccess) return 1; @@ -85,6 +83,10 @@ int main(int argc, char **argv) { if (hipStreamDestroy(Stream) != hipSuccess) return 1; + if (hipStreamDestroy(BlockingStream) != hipSuccess) + return 1; + if (hipStreamDestroy(NonBlockingStream) != hipSuccess) + return 1; print_error("destroyed stream destroy", hipStreamDestroy(Stream)); // CHECK: destroyed stream destroy value: 4 // CHECK: destroyed stream destroy name: hipErrorInvalidResourceHandle >From 0f7d7ad2acbf13ebd5319b2d603f9bed4c7e81df Mon Sep 17 00:00:00 2001 From: Sophia Herrmann <[email protected]> Date: Wed, 19 Aug 2026 17:42:04 -0700 Subject: [PATCH 2/2] add event cleanup --- offload/languages/kernel/CMakeLists.txt | 1 + .../languages/kernel/include/LanguageUtils.h | 59 ++++++++----- offload/languages/kernel/include/State.h | 3 +- offload/languages/kernel/include/Stream.h | 19 +++++ .../languages/kernel/src/LanguageLaunch.cpp | 2 + .../languages/kernel/src/LanguageRuntime.cpp | 28 +++--- offload/languages/kernel/src/State.cpp | 8 +- offload/languages/kernel/src/Stream.cpp | 85 +++++++++++++++++++ 8 files changed, 168 insertions(+), 37 deletions(-) create mode 100644 offload/languages/kernel/src/Stream.cpp diff --git a/offload/languages/kernel/CMakeLists.txt b/offload/languages/kernel/CMakeLists.txt index 52c269bed19c6..a23e5727e6b8b 100644 --- a/offload/languages/kernel/CMakeLists.txt +++ b/offload/languages/kernel/CMakeLists.txt @@ -64,6 +64,7 @@ add_llvm_library( src/LanguageLaunch.cpp src/LanguageRegistration.cpp src/State.cpp + src/Stream.cpp ) if(LLVM_LINK_LLVM_DYLIB) diff --git a/offload/languages/kernel/include/LanguageUtils.h b/offload/languages/kernel/include/LanguageUtils.h index 708de92fcdab6..62fb6160730c5 100644 --- a/offload/languages/kernel/include/LanguageUtils.h +++ b/offload/languages/kernel/include/LanguageUtils.h @@ -65,24 +65,49 @@ static inline llvm::offload::StreamTy *toInternalStream(Stream_t Stream) { return reinterpret_cast<StreamTy *>(Stream); } +static inline ol_result_t +syncAndDestroyEvents(llvm::SmallVectorImpl<ol_event_handle_t> &Events) { + ol_result_t FirstError = OL_SUCCESS; + for (ol_event_handle_t Event : Events) { + if (!Event) + continue; + + ol_result_t SyncResult = olSyncEvent(Event); + if (FirstError == OL_SUCCESS && SyncResult != OL_SUCCESS) + FirstError = SyncResult; + + ol_result_t DestroyResult = olDestroyEvent(Event); + if (FirstError == OL_SUCCESS && DestroyResult != OL_SUCCESS) + FirstError = DestroyResult; + } + Events.clear(); + return FirstError; +} + /// Wait for blocking streams before executing if we are legacy default stream. static inline ol_result_t waitOnBlockingStreams() { ol_device_handle_t Device = ThreadState::getDefaultDevice(); - if (!RuntimeState::hasLegacyDefaultStream(Device) || - RuntimeState::getBlockingStreams(Device).empty()) + llvm::SmallPtrSet<StreamTy *, 8> BlockingStreams = + RuntimeState::getBlockingStreams(Device); + if (!RuntimeState::hasLegacyDefaultStream(Device) || BlockingStreams.empty()) return OL_SUCCESS; + StreamTy *DefaultStream = ThreadState::getDefaultStream(); llvm::SmallVector<ol_event_handle_t, 8> Events; - for (StreamTy *BlockingStream : RuntimeState::getBlockingStreams(Device)) { + for (StreamTy *BlockingStream : BlockingStreams) { ol_event_handle_t Event = nullptr; ol_result_t Result = olCreateEvent(BlockingStream->Queue, OL_EVENT_FLAGS_NONE, &Event); - if (Result != OL_SUCCESS) + if (Result != OL_SUCCESS) { + if (Event) + Events.push_back(Event); + syncAndDestroyEvents(Events); return Result; + } Events.push_back(Event); } - return olWaitEvents(DefaultStream->Queue, Events.data(), Events.size()); + return DefaultStream->waitOnAndTrackDependencyEvents(Events); } /// Wait for the legacy default stream to complete before launching a kernel on @@ -99,23 +124,15 @@ static inline ol_result_t waitOnLegacyDefaultStream(StreamTy *SourceStream, ol_event_handle_t Event = nullptr; ol_result_t Result = olCreateEvent(DefaultStream->Queue, OL_EVENT_FLAGS_NONE, &Event); - if (Result != OL_SUCCESS) + if (Result != OL_SUCCESS) { + if (Event) { + llvm::SmallVector<ol_event_handle_t, 1> Events = {Event}; + syncAndDestroyEvents(Events); + } return Result; - return olWaitEvents(SourceStream->Queue, &Event, 1); -} - -/// Convert a Stream_t to an ol_queue_handle_t. -static inline Error_t getQueueFromStream(Stream_t Stream, - ol_queue_handle_t *Queue) { - if (!Stream) - return ErrorInvalidValue; - - llvm::offload::StreamTy *InternalStream = toInternalStream(Stream); - if (!llvm::offload::StateTy::isStreamRegistered(InternalStream)) - return ErrorInvalidResourceHandle; - - *Queue = InternalStream->Queue; - return Success; + } + return SourceStream->waitOnAndTrackDependencyEvents( + llvm::ArrayRef<ol_event_handle_t>(&Event, 1)); } } // namespace offload diff --git a/offload/languages/kernel/include/State.h b/offload/languages/kernel/include/State.h index e7c45bc8dec88..75f0b0dd3fc58 100644 --- a/offload/languages/kernel/include/State.h +++ b/offload/languages/kernel/include/State.h @@ -148,7 +148,8 @@ struct StateTy { static llvm::SmallPtrSet<StreamTy *, 8> getBlockingStreams(ol_device_handle_t Device); - /// Return true if \p Device has an initialized (i.e. previously used) legacy default stream. + /// Return true if \p Device has an initialized (i.e. previously used) legacy + /// default stream. static bool hasLegacyDefaultStream(ol_device_handle_t Device); /// Create a stream for \p Device and register it with the process state. diff --git a/offload/languages/kernel/include/Stream.h b/offload/languages/kernel/include/Stream.h index e66fc7000281d..d3dfd5709f40e 100644 --- a/offload/languages/kernel/include/Stream.h +++ b/offload/languages/kernel/include/Stream.h @@ -10,6 +10,10 @@ #define LLVM_OFFLOAD_LANGUAGES_KERNEL_INCLUDE_STREAM_H #include "OffloadAPI.h" +#include "llvm/ADT/ArrayRef.h" +#include "llvm/ADT/SmallVector.h" +#include <cstddef> +#include <mutex> namespace llvm { namespace offload { @@ -22,9 +26,24 @@ enum class QueueKind { }; struct StreamTy { + StreamTy(ol_queue_handle_t Queue, ol_device_handle_t Device, QueueKind Kind) + : Queue(Queue), Device(Device), Kind(Kind) {} + + ol_result_t + waitOnAndTrackDependencyEvents(llvm::ArrayRef<ol_event_handle_t> Events); + ol_result_t syncStream(); + ol_queue_handle_t Queue = nullptr; ol_device_handle_t Device = nullptr; QueueKind Kind = QueueKind::ExplicitBlocking; + +private: + ol_result_t reclaimDependencyEventsLocked(); + + static constexpr size_t MaxPendingDependencyEvents = 64; + + std::mutex DependencyEventsLock; + llvm::SmallVector<ol_event_handle_t, 8> DependencyEvents; }; } // namespace offload diff --git a/offload/languages/kernel/src/LanguageLaunch.cpp b/offload/languages/kernel/src/LanguageLaunch.cpp index a305e42cac0de..38e93b3e813e4 100644 --- a/offload/languages/kernel/src/LanguageLaunch.cpp +++ b/offload/languages/kernel/src/LanguageLaunch.cpp @@ -23,6 +23,8 @@ using llvm::offload::InvalidArgumentError; using llvm::offload::InvalidConfigurationError; using llvm::offload::InvalidDeviceError; using llvm::offload::InvalidKernelError; +using llvm::offload::waitOnBlockingStreams; +using llvm::offload::waitOnLegacyDefaultStream; /// Internal kernel launch implementation ol_result_t __llvmLaunchKernelImpl(const char *KernelID, dim3 GridDim, diff --git a/offload/languages/kernel/src/LanguageRuntime.cpp b/offload/languages/kernel/src/LanguageRuntime.cpp index a57d10688b644..dade4c8aa9939 100644 --- a/offload/languages/kernel/src/LanguageRuntime.cpp +++ b/offload/languages/kernel/src/LanguageRuntime.cpp @@ -35,7 +35,7 @@ using ThreadState = llvm::offload::ThreadStateTy; using StreamTy = llvm::offload::StreamTy; using llvm::offload::convertAndSetLastError; -using llvm::offload::getQueueFromStream; +using llvm::offload::waitOnBlockingStreams; using llvm::offload::setLastError; using llvm::offload::toInternalStream; using llvm::offload::toLanguageStream; @@ -58,7 +58,8 @@ Error_t Memcpy(void *Dst, const void *Src, size_t Size, MemcpyKind Kind) { return convertAndSetLastError(Result); } ol_device_handle_t Device = ThreadState::getDefaultDevice(); - ol_queue_handle_t Queue = ThreadState::getDefaultQueue(); + StreamTy *DefaultStream = ThreadState::getDefaultStream(); + ol_queue_handle_t Queue = DefaultStream->Queue; ol_result_t Result; switch (Kind) { case MemcpyHostToHost: { @@ -89,18 +90,16 @@ Error_t Memcpy(void *Dst, const void *Src, size_t Size, MemcpyKind Kind) { if (Result != OL_SUCCESS) return convertAndSetLastError(Result); - if (!Queue) - return convertAndSetLastError(Result); - - Result = olSyncQueue(Queue); + Result = DefaultStream->syncStream(); return convertAndSetLastError(Result); } Error_t DeviceSynchronize() { // TODO: This is not correct. We likely want to pipe this through to the // plugins. - ol_queue_handle_t Queue = ThreadState::getDefaultQueue(); - ol_result_t Result = olSyncQueue(Queue); + StreamTy *DefaultStream = ThreadState::getDefaultStream(); + ol_result_t Result = + DefaultStream ? DefaultStream->syncStream() : olSyncQueue(nullptr); return convertAndSetLastError(Result); } @@ -195,11 +194,14 @@ Error_t StreamDestroy(Stream_t Stream) { } Error_t StreamSynchronize(Stream_t Stream) { - ol_queue_handle_t Queue; - Error_t Err = getQueueFromStream(Stream, &Queue); - if (Err != Success) - return setLastError(Err); - ol_result_t Result = olSyncQueue(Queue); + if (!Stream) + return setLastError(ErrorInvalidValue); + + llvm::offload::StreamTy *InternalStream = toInternalStream(Stream); + if (!llvm::offload::StateTy::isStreamRegistered(InternalStream)) + return setLastError(ErrorInvalidResourceHandle); + + ol_result_t Result = InternalStream->syncStream(); return convertAndSetLastError(Result); } diff --git a/offload/languages/kernel/src/State.cpp b/offload/languages/kernel/src/State.cpp index 6cc78330c0c84..3d882e6434df0 100644 --- a/offload/languages/kernel/src/State.cpp +++ b/offload/languages/kernel/src/State.cpp @@ -91,7 +91,7 @@ static void destroyStreamHandle(StreamTy *&Stream) { if (!Stream) return; - olSyncQueue(Stream->Queue); + (void)Stream->syncStream(); olDestroyQueue(Stream->Queue); delete Stream; Stream = nullptr; @@ -336,8 +336,12 @@ ol_result_t StateTy::destroyStream(StreamTy *Stream) { if (!isStreamRegistered(Stream)) return &InvalidStreamError; + ol_result_t Result = Stream->syncStream(); + if (Result != OL_SUCCESS) + return Result; + get().removeStream(Stream); - ol_result_t Result = olDestroyQueue(Stream->Queue); + Result = olDestroyQueue(Stream->Queue); delete Stream; return Result; } diff --git a/offload/languages/kernel/src/Stream.cpp b/offload/languages/kernel/src/Stream.cpp new file mode 100644 index 0000000000000..a6c012259146e --- /dev/null +++ b/offload/languages/kernel/src/Stream.cpp @@ -0,0 +1,85 @@ +//===-- Stream.cpp - Kernel language stream state -------------------------===// +// +// Part of the LLVM Project, under the Apache License v2.0 with LLVM Exceptions. +// See https://llvm.org/LICENSE.txt for license information. +// SPDX-License-Identifier: Apache-2.0 WITH LLVM-exception +// +//===----------------------------------------------------------------------===// + +#include "Stream.h" +#include "OffloadAPI.h" +#include "llvm/ADT/ArrayRef.h" +#include "llvm/ADT/SmallVector.h" + +#include <mutex> +#include <utility> + +using namespace llvm; +using namespace offload; + +static ol_result_t syncAndDestroyEvents(ArrayRef<ol_event_handle_t> Events) { + ol_result_t FirstError = OL_SUCCESS; + for (ol_event_handle_t Event : Events) { + if (!Event) + continue; + + ol_result_t SyncResult = olSyncEvent(Event); + if (FirstError == OL_SUCCESS && SyncResult != OL_SUCCESS) + FirstError = SyncResult; + + ol_result_t DestroyResult = olDestroyEvent(Event); + if (FirstError == OL_SUCCESS && DestroyResult != OL_SUCCESS) + FirstError = DestroyResult; + } + return FirstError; +} + +ol_result_t +StreamTy::waitOnAndTrackDependencyEvents(ArrayRef<ol_event_handle_t> Events) { + if (Events.empty()) + return OL_SUCCESS; + + SmallVector<ol_event_handle_t, 8> MutableEvents(Events.begin(), Events.end()); + std::lock_guard<std::mutex> LG(DependencyEventsLock); + ol_result_t WaitResult = + olWaitEvents(Queue, MutableEvents.data(), MutableEvents.size()); + if (WaitResult != OL_SUCCESS) { + syncAndDestroyEvents(MutableEvents); + return WaitResult; + } + + DependencyEvents.append(MutableEvents.begin(), MutableEvents.end()); + ol_result_t ReclaimResult = OL_SUCCESS; + if (DependencyEvents.size() >= MaxPendingDependencyEvents) { + ReclaimResult = olSyncQueue(Queue); + if (ReclaimResult == OL_SUCCESS) + ReclaimResult = reclaimDependencyEventsLocked(); + } + return ReclaimResult; +} + +ol_result_t StreamTy::syncStream() { + std::lock_guard<std::mutex> LG(DependencyEventsLock); + ol_result_t Result = olSyncQueue(Queue); + if (Result != OL_SUCCESS) + return Result; + return reclaimDependencyEventsLocked(); +} + +ol_result_t StreamTy::reclaimDependencyEventsLocked() { + if (DependencyEvents.empty()) + return OL_SUCCESS; + + ol_result_t FirstError = OL_SUCCESS; + SmallVector<ol_event_handle_t, 8> RemainingEvents; + for (ol_event_handle_t Event : DependencyEvents) { + ol_result_t Result = olDestroyEvent(Event); + if (Result != OL_SUCCESS) { + if (FirstError == OL_SUCCESS) + FirstError = Result; + RemainingEvents.push_back(Event); + } + } + DependencyEvents = std::move(RemainingEvents); + return FirstError; +} _______________________________________________ llvm-branch-commits mailing list [email protected] https://lists.llvm.org/cgi-bin/mailman/listinfo/llvm-branch-commits
