[llvm-branch-commits] [clang] [llvm] [Offload][Lang] Add blocking to LaunchKernel and Memcpy (PR #218049)

Sophia Herrmann via llvm-branch-commits llvm-branch-commits at lists.llvm.org
Mon Aug 24 13:04:28 PDT 2026


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 <herrmann15 at llnl.gov>
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 <herrmann15 at llnl.gov>
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;
+}



More information about the llvm-branch-commits mailing list