From 95b0992ca0249e4eb37e056c9b96afbfb2ddc71c Mon Sep 17 00:00:00 2001 From: Tim Besard Date: Wed, 30 Sep 2026 19:39:06 +0200 Subject: [PATCH 1/3] Make synchronization cooperative, using GPUToolbox. `synchronize` blocked the calling thread in the driver until the work had completed, so no other task could run on it in the meantime, and no other thread could run the GC either. KernelInterface requires `synchronize` to be cooperative. Use GPUToolbox's `cooperative_wait`, as CUDA.jl does: poll a non-blocking query first, which keeps the latency of short operations low, and then block in the driver (GC-safe) on a separate thread, while the calling task waits without blocking the scheduler. This applies to `synchronize()` and `synchronize(::oneStream)`, i.e., to user code, KernelAbstractions and the synchronizing copies; `synchronize(; blocking=true)` restores the old behavior. The command list and queue methods, which also run from finalizers, keep blocking. With `ONEAPI_SYNC_EACH_SUBMISSION` set, as on Aurora's LTS stack, the wait after every launch is cooperative too; otherwise it would block the thread, and leave `synchronize` nothing to wait for. --- Project.toml | 2 +- lib/level-zero/oneL0.jl | 1 + lib/level-zero/synchronization.jl | 51 +++++++++++++++++++++++++++++++ lib/utils/APIUtils.jl | 4 +-- src/compiler/execution.jl | 3 +- src/context.jl | 24 +++++++++------ src/oneAPIKernels.jl | 1 - test/execution.jl | 51 +++++++++++++++++++++++++++++++ 8 files changed, 122 insertions(+), 15 deletions(-) create mode 100644 lib/level-zero/synchronization.jl diff --git a/Project.toml b/Project.toml index 88b57ba4..4dcc09b2 100644 --- a/Project.toml +++ b/Project.toml @@ -40,7 +40,7 @@ CEnum = "0.4, 0.5" ExprTools = "0.1" GPUArrays = "11.5.14" GPUCompiler = "2.9" -GPUToolbox = "3.1" +GPUToolbox = "3.2" KernelAbstractions = "0.9.39" LLVM = "6, 7, 8, 9" NEO_jll = "=26.18.38308" diff --git a/lib/level-zero/oneL0.jl b/lib/level-zero/oneL0.jl index 82e33eec..c36b9eee 100644 --- a/lib/level-zero/oneL0.jl +++ b/lib/level-zero/oneL0.jl @@ -124,6 +124,7 @@ end include("context.jl") include("cmdqueue.jl") include("cmdlist.jl") +include("synchronization.jl") include("fence.jl") include("event.jl") include("barrier.jl") diff --git a/lib/level-zero/synchronization.jl b/lib/level-zero/synchronization.jl new file mode 100644 index 00000000..e3a3a4e2 --- /dev/null +++ b/lib/level-zero/synchronization.jl @@ -0,0 +1,51 @@ +# cooperative synchronization +# +# `zeCommandListHostSynchronize` and `zeCommandQueueSynchronize` block the calling thread +# until the work has completed, so no other task can run on it in the meantime. Instead, wait +# using GPUToolbox's `cooperative_wait`: first poll, which keeps the latency of short +# operations low, and then block in the driver on a separate thread, while the calling task +# yields. + +using GPUToolbox: cooperative_wait + +export nonblocking_synchronize + +const SyncObject = Union{ZeImmediateCommandList, ZeCommandQueue} + +# with a zero timeout, a synchronization is a query +function check_done(res::ze_result_t) + if res == RESULT_NOT_READY + return false + elseif res == RESULT_SUCCESS + return true + else + throw_api_error(res) + end +end +Base.isdone(list::ZeImmediateCommandList) = + check_done(unchecked_zeCommandListHostSynchronize(list, 0)) +Base.isdone(queue::ZeCommandQueue) = + check_done(unchecked_zeCommandQueueSynchronize(queue, 0)) + +# the blocking synchronization, marked GC-safe so that it doesn't keep the GC from running +gcsafe_synchronize(list::ZeImmediateCommandList) = + @gcsafe_ccall libze_loader.zeCommandListHostSynchronize( + list::ze_command_list_handle_t, typemax(UInt64)::UInt64)::ze_result_t +gcsafe_synchronize(queue::ZeCommandQueue) = + @gcsafe_ccall libze_loader.zeCommandQueueSynchronize( + queue::ze_command_queue_handle_t, typemax(UInt64)::UInt64)::ze_result_t + +""" + nonblocking_synchronize(list_or_queue) + +Wait for the work on an immediate command list or command queue to complete, like +[`synchronize`](@ref), but without blocking the calling thread: other tasks keep running +while this one waits. +""" +function nonblocking_synchronize(obj::SyncObject) + # when polling found the work to be done, synchronize again to check for errors + res = @something(cooperative_wait(gcsafe_synchronize, obj; isdone=Base.isdone), + gcsafe_synchronize(obj)) + res == RESULT_SUCCESS || throw_api_error(res) + return +end diff --git a/lib/utils/APIUtils.jl b/lib/utils/APIUtils.jl index d8d30394..61ec0743 100644 --- a/lib/utils/APIUtils.jl +++ b/lib/utils/APIUtils.jl @@ -1,8 +1,8 @@ module APIUtils # helpers that facilitate working with C APIs -using GPUToolbox: @checked, @debug_ccall -export @checked, @debug_ccall +using GPUToolbox: @checked, @debug_ccall, @gcsafe_ccall +export @checked, @debug_ccall, @gcsafe_ccall include("enum.jl") end diff --git a/src/compiler/execution.jl b/src/compiler/execution.jl index 20742c2d..f129919b 100644 --- a/src/compiler/execution.jl +++ b/src/compiler/execution.jl @@ -360,7 +360,8 @@ end spill > s.scratch_hwm && scratch_hedge!(s, spill) append_launch!(s.list, kernel, groups) - oneL0.sync_each_submission() && oneL0.synchronize(s.list) + # wait cooperatively, as `synchronize` does, or the workaround blocks the thread + oneL0.sync_each_submission() && oneL0.nonblocking_synchronize(s.list) return end diff --git a/src/context.jl b/src/context.jl index 409a9d7a..6c9e188f 100644 --- a/src/context.jl +++ b/src/context.jl @@ -378,12 +378,15 @@ function synchronize_all_streams(ctx::ZeContext, dev::Union{ZeDevice, Nothing}) end """ - synchronize() - synchronize(stream::oneStream) + synchronize(; blocking=false) + synchronize(stream::oneStream; blocking=false) -Block the host thread until all operations on the calling task's stream for the current -context and device have completed: work appended to the immediate command list as well -as oneMKL work on the companion queue. +Block the calling task until all operations on its stream for the current context and +device have completed: work appended to the immediate command list as well as oneMKL work +on the companion queue. + +Unless `blocking` is set, other tasks keep running while waiting: the host thread is only +blocked in the driver when the work is already done. This is useful for timing operations or ensuring that GPU work has finished before accessing results on the CPU. @@ -398,18 +401,19 @@ println("GPU work completed") See also: [`global_stream`](@ref), [`context`](@ref), [`device`](@ref) """ -function oneL0.synchronize(s::oneStream) - oneL0.synchronize(s.list) +function oneL0.synchronize(s::oneStream; blocking::Bool=false) + sync = blocking ? oneL0.synchronize : oneL0.nonblocking_synchronize + sync(s.list) q = s.queue if q !== nothing - oneL0.synchronize(q) + sync(q) s.mkl_dirty = false end return end -function oneL0.synchronize() - oneL0.synchronize(global_stream(context(), device())) +function oneL0.synchronize(; blocking::Bool=false) + oneL0.synchronize(global_stream(context(), device()); blocking) end # Julia → MKL ordering: everything Julia appended to the task's immediate list must be diff --git a/src/oneAPIKernels.jl b/src/oneAPIKernels.jl index 7b90d2ca..ae693caa 100644 --- a/src/oneAPIKernels.jl +++ b/src/oneAPIKernels.jl @@ -26,7 +26,6 @@ oneAPIBackend(; prefer_blocks = false, always_inline = false) = oneAPIBackend(pr @inline KA.ones(::oneAPIBackend, ::Type{T}, dims::Tuple; unified::Bool = false) where {T} = fill!(oneArray{T, length(dims), unified ? oneAPI.oneL0.SharedBuffer : oneAPI.oneL0.DeviceBuffer}(undef, dims), one(T)) KA.get_backend(::oneArray) = oneAPIBackend() -# TODO should be non-blocking KA.synchronize(::oneAPIBackend) = oneAPI.oneL0.synchronize() KA.supports_float64(::oneAPIBackend) = false # TODO: Check if this is device dependent KA.supports_unified(::oneAPIBackend) = true diff --git a/test/execution.jl b/test/execution.jl index 3efcdbf1..c1ab7e1b 100644 --- a/test/execution.jl +++ b/test/execution.jl @@ -743,6 +743,57 @@ end @test all(results) end +# burns `iters` dependent steps per work-item, which the compiler cannot fold away +function slow_kernel(a, iters) + i = get_global_id() + acc = i % UInt32 + for k in UInt32(1):iters + acc = acc * 0x0019660d + k + end + @inbounds a[i] = acc + return +end + +@testset "cooperative synchronize" begin + a = oneArray{UInt32}(undef, 64) + slow(iters) = @oneapi items=64 slow_kernel(a, UInt32(iters)) + slow(1) + synchronize() + # warm up the slow path of `synchronize`, which would otherwise be compiled while waiting + slow(2^16) + synchronize() + + # make the kernel run for a while (calibrated with blocking synchronization, which + # does not involve other threads) + iters = 2^16 + while iters < 2^30 && @elapsed((slow(iters); synchronize(; blocking=true))) < 0.1 + iters *= 4 + end + + # other tasks keep running while one waits (when blocking, the ticker could not run + # at all until the kernel had finished) + stamps = UInt64[] + waiting = Ref(true) + ticker = @async while waiting[] + push!(stamps, time_ns()) + sleep(0.001) + end + local done + try + slow(iters) + synchronize() + done = time_ns() + finally + waiting[] = false + wait(ticker) + end + @test count(<(done), stamps) > 2 + + # blocking synchronization is still available + slow(1) + @test synchronize(; blocking=true) === nothing +end + ############################################################################################ # Keep allocation consumers at top level so kernels do not capture test state. From 72a93545848b6ad1e5806ab86241d9ed821af5ea Mon Sep 17 00:00:00 2001 From: Tim Besard Date: Thu, 1 Oct 2026 12:00:41 +0200 Subject: [PATCH 2/3] Require GPUToolbox 3.3.1 3.3.1 keeps cooperative_wait's worker threads from running finalizers. On the LTS stack, freeing memory from a finalizer drains every stream, so a worker could block in there before waking the task it waited for. --- Project.toml | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/Project.toml b/Project.toml index 4dcc09b2..2a992247 100644 --- a/Project.toml +++ b/Project.toml @@ -40,7 +40,7 @@ CEnum = "0.4, 0.5" ExprTools = "0.1" GPUArrays = "11.5.14" GPUCompiler = "2.9" -GPUToolbox = "3.2" +GPUToolbox = "3.3.1" KernelAbstractions = "0.9.39" LLVM = "6, 7, 8, 9" NEO_jll = "=26.18.38308" From de2d8c1cb4b57689d236843ed1cec7dc04acb3d7 Mon Sep 17 00:00:00 2001 From: Tim Besard Date: Thu, 1 Oct 2026 16:23:19 +0200 Subject: [PATCH 3/3] Make the cooperative synchronization test independent of timing. The test counted how often a timer task ran while a calibrated kernel kept the GPU busy, which isn't reliable on a loaded machine. Instead, keep the work on the stream waiting for a host-signalled event, which another task on the same thread only signals after yielding many times, and check that synchronizing returned after it did. A watchdog on one of cooperative_wait's worker threads signals the event after a timeout and records that, so blocking synchronization fails the test instead of hanging it. Also cover synchronizing a stream, and the wait after every launch with ONEAPI_SYNC_EACH_SUBMISSION. --- test/execution.jl | 124 ++++++++++++++++++++++++++++++++-------------- 1 file changed, 88 insertions(+), 36 deletions(-) diff --git a/test/execution.jl b/test/execution.jl index c1ab7e1b..c32cb176 100644 --- a/test/execution.jl +++ b/test/execution.jl @@ -743,55 +743,107 @@ end @test all(results) end -# burns `iters` dependent steps per work-item, which the compiler cannot fold away -function slow_kernel(a, iters) - i = get_global_id() - acc = i % UInt32 - for k in UInt32(1):iters - acc = acc * 0x0019660d + k - end - @inbounds a[i] = acc +function fill_kernel(a, x) + @inbounds a[get_global_id()] = x return end @testset "cooperative synchronize" begin - a = oneArray{UInt32}(undef, 64) - slow(iters) = @oneapi items=64 slow_kernel(a, UInt32(iters)) - slow(1) - synchronize() - # warm up the slow path of `synchronize`, which would otherwise be compiled while waiting - slow(2^16) + # keep the work on the task's stream waiting for an event until another task signals + # it, which makes these tests independent of timing: synchronizing can only return after + # that task got to run. if it cannot (e.g., because synchronizing blocks the thread), a + # watchdog on another thread signals the event after a while instead, and records that. + pool = oneL0.ZeEventPool(context(), 1, device(); + flags=oneL0.ZE_EVENT_POOL_FLAG_HOST_VISIBLE) + gate = pool[1] + timeout = UInt64(60_000_000_000) # ns + + # run `f` with the task's stream waiting for the gate, which another task on the same + # thread opens, 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, and the watchdog did not have to. + function gated(f) + list = oneAPI.global_stream(context(), device()).list + reset(gate) + # while the gate is closed, nothing on the host may wait for the GPU to become + # idle, as freeing memory does on some stacks. so avoid running finalizers, by + # collecting beforehand and not collecting while the gate is closed. + GC.gc(true) + gc_enabled = GC.enable(false) + watching = Threads.Atomic{Bool}(false) + opened = Threads.Atomic{Bool}(false) + in_time, timed_out = false, true + local watchdog, opener + try + # the watchdog waits on one of `cooperative_wait`'s worker threads, which keep + # running when this thread is blocked + watchdog = @async oneL0.cooperative_wait(gate; spin=false) do gate + watching[] = true + res = oneL0.@gcsafe_ccall oneL0.libze_loader.zeEventHostSynchronize( + gate::oneL0.ze_event_handle_t, timeout::UInt64)::oneL0.ze_result_t + res == oneL0.RESULT_NOT_READY || return false + oneL0.signal(gate) + return true + end + while !watching[] + yield() + end + + oneL0.append_wait!(list, gate) + opener = @async begin + for _ in 1:10_000 + yield() + end + opened[] = true + oneL0.signal(gate) + end + f() + in_time = opened[] + finally + # also when `f` failed, as the gate is reused + oneL0.signal(gate) + @isdefined(opener) && wait(opener) + @isdefined(watchdog) && (timed_out = something(fetch(watchdog))) + GC.enable(gc_enabled) + end + synchronize() + return in_time && !timed_out + end + + # compile everything beforehand, as that could wait for the GPU to become idle + a = oneArray{Int32}(undef, 64) + @oneapi items=64 fill_kernel(a, Int32(0)) synchronize() + @test gated(synchronize) - # make the kernel run for a while (calibrated with blocking synchronization, which - # does not involve other threads) - iters = 2^16 - while iters < 2^30 && @elapsed((slow(iters); synchronize(; blocking=true))) < 0.1 - iters *= 4 + @test gated() do + oneL0.sync_each_submission(false) do + @oneapi items=64 fill_kernel(a, Int32(1)) + end + synchronize() end + @test Array(a) == fill(Int32(1), 64) - # other tasks keep running while one waits (when blocking, the ticker could not run - # at all until the kernel had finished) - stamps = UInt64[] - waiting = Ref(true) - ticker = @async while waiting[] - push!(stamps, time_ns()) - sleep(0.001) + @test gated() do + oneL0.sync_each_submission(false) do + @oneapi items=64 fill_kernel(a, Int32(2)) + end + synchronize(oneAPI.global_stream(context(), device())) end - local done - try - slow(iters) - synchronize() - done = time_ns() - finally - waiting[] = false - wait(ticker) + @test Array(a) == fill(Int32(2), 64) + + # synchronizing after every launch, as on the LTS stack + @test gated() do + oneL0.sync_each_submission(true) do + @oneapi items=64 fill_kernel(a, Int32(3)) + end end - @test count(<(done), stamps) > 2 + @test Array(a) == fill(Int32(3), 64) # blocking synchronization is still available - slow(1) + @oneapi items=64 fill_kernel(a, Int32(4)) @test synchronize(; blocking=true) === nothing + @test Array(a) == fill(Int32(4), 64) end ############################################################################################