From f310beb9e4dbbf43ae890172cf496993c4f84b21 Mon Sep 17 00:00:00 2001 From: Tim Besard Date: Wed, 30 Sep 2026 21:50:01 +0200 Subject: [PATCH 1/5] Use GPUToolbox's cooperative_wait to synchronize streams and events. Nonblocking stream synchronization enqueued a host function that woke the waiting task through libuv, plus a task polling the stream every second in case the host function never ran. Instead, wait for the stream in HIP on a worker thread, as CUDA.jl, oneAPI.jl and OpenCL.jl do with the shared implementation in GPUToolbox. This avoids the host function, which delays both the wakeup and later work on the stream, and the extra tasks and timer per synchronization. Event synchronization used to block the calling thread after polling, and now waits the same way. --- Project.toml | 2 +- docs/src/api/streams.md | 11 +++-- src/hip/HIP.jl | 2 +- src/hip/event.jl | 35 +++++++-------- src/hip/stream.jl | 92 ++++++++++---------------------------- test/hip_core_tests.jl | 97 +++++++++++++++++++++++++++++++++++++++++ 6 files changed, 145 insertions(+), 94 deletions(-) diff --git a/Project.toml b/Project.toml index 8bc9180ef..11dac5293 100644 --- a/Project.toml +++ b/Project.toml @@ -61,7 +61,7 @@ EnzymeCore = "0.8" ExprTools = "0.1" GPUArrays = "11.5.14" GPUCompiler = "2.9" -GPUToolbox = "3" +GPUToolbox = "3.2" KernelAbstractions = "0.9.2" LLVM = "9" LLVMDowngrader_jll = "0.11" diff --git a/docs/src/api/streams.md b/docs/src/api/streams.md index 3133ccab1..ef0e6c5d6 100644 --- a/docs/src/api/streams.md +++ b/docs/src/api/streams.md @@ -56,14 +56,19 @@ AMDGPU.HIPStream ## Synchronization -AMDGPU.jl by default uses non-blocking stream synchronization with -[`AMDGPU.synchronize`](@ref) to work correctly with TLS and [Hostcall](@ref). +By default, [`AMDGPU.synchronize`](@ref) does not block the calling thread: +it briefly polls the stream, and then waits for it on a separate worker +thread, so that other tasks can run on the calling thread in the meantime. +This is required for [Hostcall](@ref), whose host side runs as a task. +Synchronizing events works the same way. +Inside finalizers, which cannot switch tasks, synchronization blocks. Users, however, can switch to a blocking synchronization globally with `nonblocking_synchronization` [preference](https://github.com/JuliaPackaging/Preferences.jl) or with fine-grained `AMDGPU.synchronize(; blocking=true)`. -Blocking synchronization might offer slightly lower latency. +Blocking synchronization might offer slightly lower latency, +but must not be used while hostcalls are running. You can also perform synchronization of the expression with [`AMDGPU.@sync`](@ref) macro, which will execute given expression and diff --git a/src/hip/HIP.jl b/src/hip/HIP.jl index 58309497a..0c3cf7c96 100644 --- a/src/hip/HIP.jl +++ b/src/hip/HIP.jl @@ -11,7 +11,7 @@ import ..AMDGPU import ..AMDGPU.libhip import .AMDGPU: @check, check -import GPUToolbox: @gcsafe_ccall, @checked +import GPUToolbox: @gcsafe_ccall, @checked, cooperative_wait include("libhip.jl") include("error.jl") diff --git a/src/hip/event.jl b/src/hip/event.jl index cb39c2dc9..ca27e5b77 100644 --- a/src/hip/event.jl +++ b/src/hip/event.jl @@ -22,30 +22,25 @@ function isdone(event::HIPEvent) end end -function non_blocking_synchronize(event::HIPEvent) - isdone(event) && return true +wait(event::HIPEvent) = hipEventSynchronize(event) + +# same, but callable from any thread (events know their device) +worker_synchronize(event::HIPEvent) = + @gcsafe_ccall(libhip.hipEventSynchronize(event::hipEvent_t)::hipError_t) - # spin (initially without yielding to minimize latency) - spins = 0 - while spins < 256 - if spins < 32 - ccall(:jl_cpu_pause, Cvoid, ()) - # Temporary solution before we have gc transition support in codegen. - ccall(:jl_gc_safepoint, Cvoid, ()) +function synchronize(event::HIPEvent; blocking::Bool = false, spin::Bool = true) + if use_nonblocking_synchronize && !blocking + res = cooperative_wait(worker_synchronize, event; isdone, spin) + if res === nothing + wait(event) else - yield() + check(something(res)) + AMDGPU.maybe_collect(; blocking=true) end - isdone(event) && return true - spins += 1 + else + AMDGPU.maybe_collect(; blocking=true) + wait(event) end - return false -end - -wait(event::HIPEvent) = hipEventSynchronize(event) - -function synchronize(event::HIPEvent) - non_blocking_synchronize(event) || AMDGPU.maybe_collect(; blocking=true) - wait(event) return end diff --git a/src/hip/stream.jl b/src/hip/stream.jl index b2c1dbc9b..70fce852a 100644 --- a/src/hip/stream.jl +++ b/src/hip/stream.jl @@ -62,83 +62,37 @@ function isdone(stream::HIPStream) end end -function _low_latency_synchronize(stream::HIPStream) - isdone(stream) && return true - - # spin (initially without yielding to minimize latency) - spins = 0 - while spins < 256 - if spins < 32 - ccall(:jl_cpu_pause, Cvoid, ()) - # Temporary solution before we have gc transition support in codegen. - ccall(:jl_gc_safepoint, Cvoid, ()) - else - yield() - end - isdone(stream) && return true - spins += 1 - end - return false -end +wait(stream::HIPStream) = hipStreamSynchronize(stream) -function launch(f::Base.Callable; stream::HIPStream) - # Condition object is embedded in a task, Julia scheduler keeps it alive. - cond = Base.AsyncCondition() do async_cond - f() - close(async_cond) - end - callback = cglobal(:uv_async_send) - hipLaunchHostFunc(stream, callback, cond) +# same, but callable from any thread. this bypasses the task-local state, so select the +# caller's device ourselves: the null stream refers to the current device's. +function worker_synchronize(stream::HIPStream, dev::HIPDevice) + isvalid(stream) || return hipSuccess + res = unchecked_hipSetDevice(device_id(dev)) + res == hipSuccess || return res + @gcsafe_ccall(libhip.hipStreamSynchronize(stream::hipStream_t)::hipError_t) end -function nonblocking_synchronize(stream::HIPStream) - # Wait for an event signalled by HIP. - event = Base.Event() - launch(() -> notify(event); stream) - - # If an error occurs, the callback may never fire. - # Create a timer to detect such cases. - dev = device() - timer = Timer(0; interval=1) - - Base.@sync begin - # Launch timer. - Threads.@spawn try - device!(dev) - while true - try - Base.wait(timer) - catch err - err isa EOFError && break - rethrow() - end - (!isvalid(stream) || hipStreamQuery(stream) != hipErrorNotReady) && break - end - finally - notify(event) - end - # Wait for `event`. - Threads.@spawn begin - Base.wait(event) - close(timer) - end - end - return -end - -wait(stream::HIPStream) = hipStreamSynchronize(stream) - -function synchronize(stream::HIPStream; blocking::Bool = false) - if use_nonblocking_synchronize && !blocking - if !_low_latency_synchronize(stream) - nonblocking_synchronize(stream) +function synchronize(stream::HIPStream; blocking::Bool = false, spin::Bool = true) + if GC.in_finalizer() + # we can't switch tasks here, and the finalizer selected the context to use + wait(stream) + elseif use_nonblocking_synchronize && !blocking + # wait on a worker thread, so that other tasks (e.g. hostcalls) can run on this one + dev = AMDGPU.device() + res = cooperative_wait(s -> worker_synchronize(s, dev), stream; isdone, spin) + if res === nothing + # polling found the stream done. synchronize anyway, which reports errors and + # lets HIP release resources. + wait(stream) + else + check(something(res)) AMDGPU.maybe_collect(; blocking=true) end else AMDGPU.maybe_collect(; blocking=true) + wait(stream) end - # Perform an actual API call even after non-blocking synchronization. - wait(stream) return end diff --git a/test/hip_core_tests.jl b/test/hip_core_tests.jl index 8e38241a3..faae25cf8 100644 --- a/test/hip_core_tests.jl +++ b/test/hip_core_tests.jl @@ -22,6 +22,103 @@ Random.seed!(1) @test t >= 0 end +@testset "cooperative synchronization" begin + # keep the GPU busy for a while + function sleep_kernel(n) + for _ in 1:n + AMDGPU.Device.device_sleep(Int32(127)) + end + return + end + function busy(n; stream) + @roc stream=stream sleep_kernel(n) + return + end + + # run `f` while counting how often another task on the same thread gets to run. polling + # before waiting yields a couple of hundred times, so only much larger counts show that + # the thread was not blocked while waiting. + function progress_during(f) + progress = Ref(0) + done = Ref(false) + t = @async while !done[] + progress[] += 1 + yield() + end + try + f() + finally + done[] = true + wait(t) + end + return progress[] + end + + # warm up everything that is measured below + let s = HIPStream() + busy(1; stream=s) + progress_during(() -> AMDGPU.synchronize(s)) + progress_during(() -> HIP.synchronize(HIP.HIPEvent(s))) + progress_during(() -> AMDGPU.synchronize(s; blocking=true)) + end + + # find a kernel that takes at least 200 ms + n = 1000 + while true + s = HIPStream() + t = @elapsed (busy(n; stream=s); AMDGPU.synchronize(s; blocking=true)) + t >= 0.2 && break + n *= 2 + end + + let s = HIPStream() + busy(n; stream=s) + @test !HIP.isdone(s) + @test progress_during(() -> AMDGPU.synchronize(s)) > 1000 + @test HIP.isdone(s) + end + + let s = HIPStream() + busy(n; stream=s) + @test !HIP.isdone(s) + @test progress_during(() -> HIP.synchronize(s; spin=false)) > 1000 + end + + let s = HIPStream() + busy(n; stream=s) + e = HIP.HIPEvent(s) + @test !HIP.isdone(e) + @test progress_during(() -> HIP.synchronize(e)) > 1000 + @test HIP.isdone(e) + end + + # the null stream belongs to the current device, which the worker has to select + let s = HIP.default_stream() + busy(n; stream=s) + @test !HIP.isdone(s) + @test progress_during(() -> AMDGPU.synchronize(s)) > 1000 + @test HIP.isdone(s) + end + + # opting out + let s = HIPStream() + busy(n; stream=s) + @test !HIP.isdone(s) + @test progress_during(() -> AMDGPU.synchronize(s; blocking=true)) < 1000 + end + + if length(AMDGPU.devices()) > 1 + # waiting for another device doesn't change the one this task uses + dev = AMDGPU.device() + other = first(d for d in AMDGPU.devices() if d != dev) + s = AMDGPU.device!(() -> HIPStream(), other) + AMDGPU.device!(() -> busy(1; stream=s), other) + HIP.synchronize(s; spin=false) + @test HIP.isdone(s) + @test AMDGPU.device() == dev + end +end + if length(AMDGPU.devices()) > 1 @testset "HIP Peer Access" begin dev1, dev2 = AMDGPU.devices()[1:2] From b07213430554a390d5693f7067446b3f239cc06b Mon Sep 17 00:00:00 2001 From: Tim Besard Date: Wed, 30 Sep 2026 21:59:37 +0200 Subject: [PATCH 2/5] Don't block the thread in device_synchronize. HIP.device_synchronize called hipDeviceSynchronize on the calling thread, which deadlocked when a kernel was waiting for a hostcall whose task could only run on that thread. Wait on a worker thread instead, like stream and event synchronization do. --- docs/src/api/streams.md | 2 +- src/hip/HIP.jl | 16 ++++++++++++++-- test/device/hostcall.jl | 33 +++++++++++++++++++++++++++++++++ test/hip_core_tests.jl | 9 +++++++++ 4 files changed, 57 insertions(+), 3 deletions(-) diff --git a/docs/src/api/streams.md b/docs/src/api/streams.md index ef0e6c5d6..9b8ea4a39 100644 --- a/docs/src/api/streams.md +++ b/docs/src/api/streams.md @@ -60,7 +60,7 @@ By default, [`AMDGPU.synchronize`](@ref) does not block the calling thread: it briefly polls the stream, and then waits for it on a separate worker thread, so that other tasks can run on the calling thread in the meantime. This is required for [Hostcall](@ref), whose host side runs as a task. -Synchronizing events works the same way. +Synchronizing events and `HIP.device_synchronize()` works the same way. Inside finalizers, which cannot switch tasks, synchronization blocks. Users, however, can switch to a blocking synchronization globally diff --git a/src/hip/HIP.jl b/src/hip/HIP.jl index 0c3cf7c96..87c417b7d 100644 --- a/src/hip/HIP.jl +++ b/src/hip/HIP.jl @@ -93,12 +93,24 @@ include("pool.jl") include("module.jl") include("graph.jl") +# callable from any thread; there is no way to poll an entire device +function worker_device_synchronize(dev::HIPDevice) + res = unchecked_hipSetDevice(device_id(dev)) + res == hipSuccess || return res + @gcsafe_ccall(libhip.hipDeviceSynchronize()::hipError_t) +end + """ Blocks until all kernels on all streams have completed. Uses currently active device. """ -function device_synchronize() - hipDeviceSynchronize() +function device_synchronize(; blocking::Bool = false) + if use_nonblocking_synchronize && !blocking && !GC.in_finalizer() + res = cooperative_wait(worker_device_synchronize, AMDGPU.device()) + check(something(res)) + else + hipDeviceSynchronize() + end AMDGPU.synchronize() # To trigger any Julia-kernel exception. AMDGPU.maybe_collect(; blocking=true) return diff --git a/test/device/hostcall.jl b/test/device/hostcall.jl index b37c1a9ce..c6423e90c 100644 --- a/test/device/hostcall.jl +++ b/test/device/hostcall.jl @@ -28,6 +28,39 @@ using AMDGPU.Device: HostCallHolder, hostcall! AMDGPU.Device.free!(hc) end +@testset "Call: while waiting for the device" begin + # the host side of a hostcall is a task. with a single thread, it can only run while + # the thread that waits for the device is not blocked. + code = """ + using AMDGPU + using AMDGPU.Device: HostCallHolder, hostcall! + + function kernel(a,b,sig) + hostcall!(sig) + b[1] = a[1] + nothing + end + + RA = ROCArray(ones(Float32, 1)) + RB = ROCArray(zeros(Float32, 1)) + hc = HostCallHolder(Nothing, Tuple{}) do + nothing + end + + @roc kernel(RA, RB, hc) + AMDGPU.HIP.device_synchronize() + Array(RB)[1] == 1f0 || exit(1) + """ + cmd = `$(Base.julia_cmd()) --threads=1 --project=$(Base.active_project()) -e $code` + proc = run(pipeline(cmd; stdout, stderr); wait=false) + timer = Timer(600) do _ + kill(proc) + end + wait(proc) + close(timer) + @test success(proc) +end + @testset "Call: Error" begin function kernel(a,b,sig) hostcall!(sig) diff --git a/test/hip_core_tests.jl b/test/hip_core_tests.jl index faae25cf8..3e0a78b94 100644 --- a/test/hip_core_tests.jl +++ b/test/hip_core_tests.jl @@ -59,6 +59,7 @@ end busy(1; stream=s) progress_during(() -> AMDGPU.synchronize(s)) progress_during(() -> HIP.synchronize(HIP.HIPEvent(s))) + progress_during(HIP.device_synchronize) progress_during(() -> AMDGPU.synchronize(s; blocking=true)) end @@ -92,6 +93,13 @@ end @test HIP.isdone(e) end + let s = HIPStream() + busy(n; stream=s) + @test !HIP.isdone(s) + @test progress_during(HIP.device_synchronize) > 1000 + @test HIP.isdone(s) + end + # the null stream belongs to the current device, which the worker has to select let s = HIP.default_stream() busy(n; stream=s) @@ -115,6 +123,7 @@ end AMDGPU.device!(() -> busy(1; stream=s), other) HIP.synchronize(s; spin=false) @test HIP.isdone(s) + AMDGPU.device!(HIP.device_synchronize, other) @test AMDGPU.device() == dev end end From 863c224cca566f6e9e878ba004abe839ea51c6a4 Mon Sep 17 00:00:00 2001 From: Tim Besard Date: Thu, 1 Oct 2026 12:00:41 +0200 Subject: [PATCH 3/5] Require GPUToolbox 3.3.1. HIP.device_synchronize() waits for something that cannot be polled. With earlier versions, cooperative_wait made such waits wait for a free worker when all of them were busy, and workers could get stuck running finalizers before waking up the waiting task. Both could deadlock when the GPU work depends on the waiting task's thread making progress, as with hostcalls. --- Project.toml | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/Project.toml b/Project.toml index 11dac5293..1e3f223c8 100644 --- a/Project.toml +++ b/Project.toml @@ -61,7 +61,7 @@ EnzymeCore = "0.8" ExprTools = "0.1" GPUArrays = "11.5.14" GPUCompiler = "2.9" -GPUToolbox = "3.2" +GPUToolbox = "3.3.1" KernelAbstractions = "0.9.2" LLVM = "9" LLVMDowngrader_jll = "0.11" From f0b5b32ced00d11e9c22589a78353685b269cbf2 Mon Sep 17 00:00:00 2001 From: Tim Besard Date: Thu, 1 Oct 2026 15:40:08 +0200 Subject: [PATCH 4/5] Make the cooperative synchronization tests robust against descheduling. The tests launched a kernel and then checked that other tasks made progress while waiting for it. On a loaded node (the MI250 runner regularly runs tests 5-30x slower than usual, also on main), the process can be descheduled for longer than the kernel takes, e.g. while compiling the measurement after the launch. The kernel then completed before the wait started, which looked like a blocked thread. The wall-clock calibration was affected too, picking kernels that took only a few ms. Time the calibration kernel on the GPU, compile everything before launching the kernel, and only accept a measurement when the wait took a while, retrying with a longer kernel otherwise. --- test/hip_core_tests.jl | 60 ++++++++++++++++++++++++------------------ 1 file changed, 34 insertions(+), 26 deletions(-) diff --git a/test/hip_core_tests.jl b/test/hip_core_tests.jl index 3e0a78b94..aec386c4a 100644 --- a/test/hip_core_tests.jl +++ b/test/hip_core_tests.jl @@ -63,56 +63,64 @@ end progress_during(() -> AMDGPU.synchronize(s; blocking=true)) end - # find a kernel that takes at least 200 ms + # find a kernel that takes at least 200 ms. time it on the GPU, as this process getting + # descheduled (as happens on loaded CI nodes) would make it seem to take longer. n = 1000 - while true - s = HIPStream() - t = @elapsed (busy(n; stream=s); AMDGPU.synchronize(s; blocking=true)) - t >= 0.2 && break + while AMDGPU.@elapsed(busy(n; stream=AMDGPU.stream())) < 0.2 n *= 2 end + # measure the progress made while `sync()` waits for a kernel on `s`. that only shows + # whether the thread was blocked if the kernel kept running for a while after the wait + # started, which isn't the case when this process gets descheduled for longer than the + # kernel takes. `isdone` can't tell, as HIP may report a completed stream as busy for a + # while, so check how long the wait took instead, and if it was too short, try again + # with a longer kernel. + function progress_while_busy(sync, s) + m = n + for _ in 1:5 + busy(m; stream=s) + t = Ref(0.0) + progress = progress_during(() -> t[] = @elapsed sync()) + t[] >= 0.05 && return progress + m *= 2 + end + error("the kernel kept completing before the wait started") + end + let s = HIPStream() - busy(n; stream=s) - @test !HIP.isdone(s) - @test progress_during(() -> AMDGPU.synchronize(s)) > 1000 + @test progress_while_busy(() -> AMDGPU.synchronize(s), s) > 1000 @test HIP.isdone(s) end let s = HIPStream() - busy(n; stream=s) - @test !HIP.isdone(s) - @test progress_during(() -> HIP.synchronize(s; spin=false)) > 1000 + @test progress_while_busy(() -> HIP.synchronize(s; spin=false), s) > 1000 end - let s = HIPStream() - busy(n; stream=s) - e = HIP.HIPEvent(s) - @test !HIP.isdone(e) - @test progress_during(() -> HIP.synchronize(e)) > 1000 - @test HIP.isdone(e) + let s = HIPStream(), e = Ref{HIP.HIPEvent}() + # record the event after the kernel + sync = () -> begin + e[] = HIP.HIPEvent(s) + HIP.synchronize(e[]) + end + @test progress_while_busy(sync, s) > 1000 + @test HIP.isdone(e[]) end let s = HIPStream() - busy(n; stream=s) - @test !HIP.isdone(s) - @test progress_during(HIP.device_synchronize) > 1000 + @test progress_while_busy(HIP.device_synchronize, s) > 1000 @test HIP.isdone(s) end # the null stream belongs to the current device, which the worker has to select let s = HIP.default_stream() - busy(n; stream=s) - @test !HIP.isdone(s) - @test progress_during(() -> AMDGPU.synchronize(s)) > 1000 + @test progress_while_busy(() -> AMDGPU.synchronize(s), s) > 1000 @test HIP.isdone(s) end # opting out let s = HIPStream() - busy(n; stream=s) - @test !HIP.isdone(s) - @test progress_during(() -> AMDGPU.synchronize(s; blocking=true)) < 1000 + @test progress_while_busy(() -> AMDGPU.synchronize(s; blocking=true), s) < 1000 end if length(AMDGPU.devices()) > 1 From 1d1e639a1b52041f02c33a04b5855a6bd921c967 Mon Sep 17 00:00:00 2001 From: Tim Besard Date: Thu, 1 Oct 2026 16:21:52 +0200 Subject: [PATCH 5/5] Test cooperative synchronization without counting progress. The previous tests counted how often another task ran while waiting for a kernel, which depends on timing: when the process gets descheduled for longer than the kernel takes, as happens on the MI250 runner, the kernel completes before the wait starts, which looks like a blocked thread. Retrying with longer kernels made that less likely, but not impossible. Instead, keep the GPU busy with a kernel that spins until a flag in host memory is set, and set that flag from another task on the same thread, after it has yielded many more times than the polling at the start of a wait does. The wait can then only return if that task ran in the meantime. If it doesn't, the kernel gives up after a while and records that it timed out, so a blocking implementation fails instead of hanging. `blocking=true` is checked the same way, expecting the kernel to time out. --- test/hip_core_tests.jl | 153 ++++++++++++++++++++++------------------- 1 file changed, 82 insertions(+), 71 deletions(-) diff --git a/test/hip_core_tests.jl b/test/hip_core_tests.jl index aec386c4a..63d236ce9 100644 --- a/test/hip_core_tests.jl +++ b/test/hip_core_tests.jl @@ -23,112 +23,123 @@ Random.seed!(1) end @testset "cooperative synchronization" begin - # keep the GPU busy for a while - function sleep_kernel(n) - for _ in 1:n + # keep the GPU busy until the host opens a gate. this keeps the tests below independent + # of timing: a synchronization can only return after the task that opens the gate has + # run. if that does not happen (e.g., because the thread it runs on is blocked), the + # kernel gives up after `limit` sleeps, and records that it timed out, instead of hanging. + function gate_kernel(gate::Ptr{UInt32}, limit) + for _ in 1:limit + unsafe_load(gate, :acquire) != 0 && return AMDGPU.Device.device_sleep(Int32(127)) end + unsafe_store!(gate, UInt32(1), 2) return end - function busy(n; stream) - @roc stream=stream sleep_kernel(n) - return + gate_buf = Mem.HostBuffer(2 * sizeof(UInt32), HIP.hipHostMallocCoherent) + gate = unsafe_wrap(Array, Ptr{UInt32}(gate_buf.ptr), 2) # (is open, timed out) + gate_ptr = Ptr{UInt32}(gate_buf.dev_ptr) + # a sleep takes 127 * 64 cycles, so this takes at least 20 s at current clock rates + timeout = 7_500_000 + open_gate() = unsafe_store!(pointer(gate), UInt32(1), :release) + gate_is_open() = unsafe_load(pointer(gate), :acquire) != 0 + + # run `f` while a kernel on `stream` keeps the GPU busy until the gate is opened, + # returning what `f` returned and whether the kernel timed out. + function gated(f, stream; limit = timeout) + # while the gate is closed, nothing on the host may wait for the GPU to become idle, + # as e.g. freeing memory does. so avoid running finalizers, by collecting beforehand + # and not collecting while the gate is closed. + GC.gc(true) + gc_enabled = GC.enable(false) + ret = try + gate .= 0 + @roc stream=stream gate_kernel(gate_ptr, limit) + f() + finally + open_gate() + GC.enable(gc_enabled) + # also when `f` failed, as the gate is reused + AMDGPU.synchronize(stream) + end + return ret, gate[2] != 0 end - # run `f` while counting how often another task on the same thread gets to run. polling - # before waiting yields a couple of hundred times, so only much larger counts show that - # the thread was not blocked while waiting. - function progress_during(f) - progress = Ref(0) - done = Ref(false) - t = @async while !done[] - progress[] += 1 - yield() + # run `f` while another task on the same thread opens the gate, but only after it got + # to run many more times than the polling at the start of a synchronization yields. + # returns whether `f` only returned after the gate had been opened. + function open_gate_during(f) + t = @async begin + for _ in 1:10_000 + yield() + end + open_gate() end try f() + gate_is_open() finally - done[] = true wait(t) end - return progress[] - end - - # warm up everything that is measured below - let s = HIPStream() - busy(1; stream=s) - progress_during(() -> AMDGPU.synchronize(s)) - progress_during(() -> HIP.synchronize(HIP.HIPEvent(s))) - progress_during(HIP.device_synchronize) - progress_during(() -> AMDGPU.synchronize(s; blocking=true)) - end - - # find a kernel that takes at least 200 ms. time it on the GPU, as this process getting - # descheduled (as happens on loaded CI nodes) would make it seem to take longer. - n = 1000 - while AMDGPU.@elapsed(busy(n; stream=AMDGPU.stream())) < 0.2 - n *= 2 end - # measure the progress made while `sync()` waits for a kernel on `s`. that only shows - # whether the thread was blocked if the kernel kept running for a while after the wait - # started, which isn't the case when this process gets descheduled for longer than the - # kernel takes. `isdone` can't tell, as HIP may report a completed stream as busy for a - # while, so check how long the wait took instead, and if it was too short, try again - # with a longer kernel. - function progress_while_busy(sync, s) - m = n - for _ in 1:5 - busy(m; stream=s) - t = Ref(0.0) - progress = progress_during(() -> t[] = @elapsed sync()) - t[] >= 0.05 && return progress - m *= 2 - end - error("the kernel kept completing before the wait started") + # set up everything beforehand: compiling and loading the kernel, or creating the + # queue backing a stream (which HIP does when first using it), may wait for the GPU. + streams = [HIPStream() for _ in 1:5] + event = HIP.HIPEvent(streams[3]; do_record=false) + open_gate() + for s in (streams..., HIP.default_stream(), AMDGPU.stream()) + @roc stream=s gate_kernel(gate_ptr, 1) + AMDGPU.synchronize(s) end - let s = HIPStream() - @test progress_while_busy(() -> AMDGPU.synchronize(s), s) > 1000 - @test HIP.isdone(s) + let s = streams[1] + @test gated(s) do + open_gate_during(() -> AMDGPU.synchronize(s)) && HIP.isdone(s) + end == (true, false) end - let s = HIPStream() - @test progress_while_busy(() -> HIP.synchronize(s; spin=false), s) > 1000 + let s = streams[2] + @test gated(s) do + open_gate_during(() -> HIP.synchronize(s; spin=false)) && HIP.isdone(s) + end == (true, false) end - let s = HIPStream(), e = Ref{HIP.HIPEvent}() - # record the event after the kernel - sync = () -> begin - e[] = HIP.HIPEvent(s) - HIP.synchronize(e[]) - end - @test progress_while_busy(sync, s) > 1000 - @test HIP.isdone(e[]) + let s = streams[3] + @test gated(s) do + HIP.record(event) + open_gate_during(() -> HIP.synchronize(event)) && HIP.isdone(event) + end == (true, false) end - let s = HIPStream() - @test progress_while_busy(HIP.device_synchronize, s) > 1000 - @test HIP.isdone(s) + let s = streams[4] + @test gated(s) do + open_gate_during(HIP.device_synchronize) && HIP.isdone(s) + end == (true, false) end # the null stream belongs to the current device, which the worker has to select let s = HIP.default_stream() - @test progress_while_busy(() -> AMDGPU.synchronize(s), s) > 1000 - @test HIP.isdone(s) + @test gated(s) do + open_gate_during(() -> AMDGPU.synchronize(s)) && HIP.isdone(s) + end == (true, false) end - # opting out - let s = HIPStream() - @test progress_while_busy(() -> AMDGPU.synchronize(s; blocking=true), s) < 1000 + # opting out blocks the thread, so the gate can only open once the kernel gave up + let s = streams[5] + @test gated(s; limit = 10_000) do + open_gate_during(() -> AMDGPU.synchronize(s; blocking=true)) + end == (false, true) end + Mem.free(gate_buf) + + noop_kernel() = return if length(AMDGPU.devices()) > 1 # waiting for another device doesn't change the one this task uses dev = AMDGPU.device() other = first(d for d in AMDGPU.devices() if d != dev) s = AMDGPU.device!(() -> HIPStream(), other) - AMDGPU.device!(() -> busy(1; stream=s), other) + AMDGPU.device!(() -> (@roc stream=s noop_kernel()), other) HIP.synchronize(s; spin=false) @test HIP.isdone(s) AMDGPU.device!(HIP.device_synchronize, other)