Skip to content

Make synchronization cooperative, using GPUToolbox - #654

Merged
maleadt merged 3 commits into
mainfrom
tb/cooperative-wait
Oct 2, 2026
Merged

maleadt merged 3 commits into
mainfrom
tb/cooperative-wait

Conversation

@maleadt

@maleadt maleadt commented Sep 30, 2026 •

Copy link
Copy Markdown
Member

synchronize() currently blocks the calling thread inside the Level Zero driver until the GPU is done. While it waits, no other Julia task can run on that thread, so work like launching kernels from another task, doing I/O, or answering a heartbeat stalls behind a long-running kernel. KernelInterface also requires back-ends to synchronize cooperatively.

This PR makes it cooperative using cooperative_wait from GPUToolbox 3.3.1 (JuliaGPU/GPUToolbox.jl#23), the same implementation CUDA.jl, OpenCL.jl and KernelAbstractions' POCL back-end are switching to:

  • Short operations are detected by polling (zeCommandListHostSynchronize with a zero timeout), so they don't pay for involving another thread.
  • Longer ones are waited for in the driver on a small pool of worker threads, using a GC-safe call, while the calling task yields to others.

For example, a task that ticks every millisecond keeps ticking while another task waits for a long kernel:

ticker = @async while true
    println("tick")
    sleep(0.001)
end
@oneapi items=64 slow_kernel(a, iters)   # runs for a few hundred ms
synchronize()                            # before: no ticks until the kernel finishes

What changes:

  • synchronize() and synchronize(::oneStream), which also back KernelAbstractions' synchronize and the synchronizing copies, now wait cooperatively. synchronize(; blocking=true) gives the old behavior.
  • With ONEAPI_SYNC_EACH_SUBMISSION (the Aurora LTS workaround), the wait after every launch is cooperative too. Otherwise that wait would block the thread and leave nothing for synchronize to wait for.
  • The command list and queue methods (oneL0.synchronize(list)) and the stream drain before freeing memory on the LTS stack keep blocking, because they run from finalizers or with finalizers disabled, where a task can't yield. The internal barriers that order Julia and oneMKL work are also left as they are.
  • This requires GPUToolbox 3.3.1, whose worker threads no longer run finalizers (Don't let cooperative_wait's worker threads deadlock GPUToolbox.jl#26). On the LTS stack, freeing memory from a finalizer drains every stream, so with 3.2 a worker could block in there, draining the whole device before waking up the task it was waiting for.

The tests don't depend on timing. They put a wait on a host-signalled Level Zero event in front of the work on the task's stream, so it can't complete until the host signals that event. Then they call synchronize() (or synchronize(stream), or launch a kernel with sync-each-submission enabled) while another task on the same thread signals the event after yielding 10,000 times, well past the polling phase. Each test checks that the call returned only after that task had run, and that the kernel's result is there. If synchronization blocked the thread, the signalling task couldn't run and the call would hang. To avoid that, a watchdog on one of cooperative_wait's worker threads waits for the event with a 60 s timeout, signals it itself if it times out, and records that, so the test fails instead. I checked this by swapping blocking synchronization back in: every gated test fails after the timeout, on 1.10, 1.12 and 1.13, with 1 or 4 threads, and with and without ONEAPI_SYNC_EACH_SUBMISSION.

Cost, measured on an Iris Xe (sagittarius), alternating blocking and cooperative synchronization after a launch:

kernel blocking cooperative
~0 85.6 µs 85.6 µs
~0.7 ms 678 µs 688 µs
~2.5 ms 2521 µs 2555 µs
~10 ms 9683 µs 9748 µs

Short operations are unaffected. Longer ones pay a few tens of µs, about 1%, to wake up the waiting task. That machine was busy running other jobs, which makes the wake-ups slower than they'd normally be.

The LTS check failed once on reductions of host-accessible arrays (sum(a) returned 0 on Julia 1.13). That failure isn't caused by this PR: the same test failed the same way on Julia 1.13 before this PR existed, on #649's own branch (5c33db5) and on the unsafe_wrap branch (24ba399). Synchronization was still blocking then. Across the LTS runs that include the test, it failed in 3 of 6 runs on Julia 1.13 and never on 1.10.

`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.
@github-actions

github-actions Bot commented Sep 30, 2026 •

Copy link
Copy Markdown
Contributor

Your PR requires formatting changes to meet the project's style guidelines.
Please consider running Runic (git runic main) to apply these changes.

Click here to view the suggested changes.
diff --git a/lib/level-zero/synchronization.jl b/lib/level-zero/synchronization.jl
index e3a3a4e..f6c2aff 100644
--- a/lib/level-zero/synchronization.jl
+++ b/lib/level-zero/synchronization.jl
@@ -30,10 +30,12 @@ Base.isdone(queue::ZeCommandQueue) =
 # 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
+    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
+    queue::ze_command_queue_handle_t, typemax(UInt64)::UInt64
+)::ze_result_t
 
 """
     nonblocking_synchronize(list_or_queue)
@@ -44,8 +46,10 @@ 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 = @something(
+        cooperative_wait(gcsafe_synchronize, obj; isdone = Base.isdone),
+        gcsafe_synchronize(obj)
+    )
     res == RESULT_SUCCESS || throw_api_error(res)
     return
 end
diff --git a/src/context.jl b/src/context.jl
index 6c9e188..36f0cd0 100644
--- a/src/context.jl
+++ b/src/context.jl
@@ -401,7 +401,7 @@ println("GPU work completed")
 
 See also: [`global_stream`](@ref), [`context`](@ref), [`device`](@ref)
 """
-function oneL0.synchronize(s::oneStream; blocking::Bool=false)
+function oneL0.synchronize(s::oneStream; blocking::Bool = false)
     sync = blocking ? oneL0.synchronize : oneL0.nonblocking_synchronize
     sync(s.list)
     q = s.queue
@@ -412,8 +412,8 @@ function oneL0.synchronize(s::oneStream; blocking::Bool=false)
     return
 end
 
-function oneL0.synchronize(; blocking::Bool=false)
-    oneL0.synchronize(global_stream(context(), device()); blocking)
+function oneL0.synchronize(; blocking::Bool = false)
+    return 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/test/execution.jl b/test/execution.jl
index c32cb17..b0f232b 100644
--- a/test/execution.jl
+++ b/test/execution.jl
@@ -753,8 +753,10 @@ end
     # 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)
+    pool = oneL0.ZeEventPool(
+        context(), 1, device();
+        flags = oneL0.ZE_EVENT_POOL_FLAG_HOST_VISIBLE
+    )
     gate = pool[1]
     timeout = UInt64(60_000_000_000)    # ns
 
@@ -777,10 +779,11 @@ end
         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
+            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
+                    gate::oneL0.ze_event_handle_t, timeout::UInt64
+                )::oneL0.ze_result_t
                 res == oneL0.RESULT_NOT_READY || return false
                 oneL0.signal(gate)
                 return true
@@ -812,13 +815,13 @@ 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))
+    @oneapi items = 64 fill_kernel(a, Int32(0))
     synchronize()
     @test gated(synchronize)
 
     @test gated() do
         oneL0.sync_each_submission(false) do
-            @oneapi items=64 fill_kernel(a, Int32(1))
+            @oneapi items = 64 fill_kernel(a, Int32(1))
         end
         synchronize()
     end
@@ -826,7 +829,7 @@ end
 
     @test gated() do
         oneL0.sync_each_submission(false) do
-            @oneapi items=64 fill_kernel(a, Int32(2))
+            @oneapi items = 64 fill_kernel(a, Int32(2))
         end
         synchronize(oneAPI.global_stream(context(), device()))
     end
@@ -835,14 +838,14 @@ end
     # 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))
+            @oneapi items = 64 fill_kernel(a, Int32(3))
         end
     end
     @test Array(a) == fill(Int32(3), 64)
 
     # blocking synchronization is still available
-    @oneapi items=64 fill_kernel(a, Int32(4))
-    @test synchronize(; blocking=true) === nothing
+    @oneapi items = 64 fill_kernel(a, Int32(4))
+    @test synchronize(; blocking = true) === nothing
     @test Array(a) == fill(Int32(4), 64)
 end
 

@codecov

codecov Bot commented Oct 1, 2026 •

Copy link
Copy Markdown

Codecov Report

❌ Patch coverage is 95.23810% with 1 line in your changes missing coverage. Please review.
✅ Project coverage is 80.90%. Comparing base (5056aa9) to head (de2d8c1).

Files with missing lines Patch % Lines
lib/level-zero/synchronization.jl 92.85% 1 Missing ⚠️
Additional details and impacted files
@@            Coverage Diff             @@
##             main     #654      +/-   ##
==========================================
+ Coverage   80.88%   80.90%   +0.02%     
==========================================
  Files          56       57       +1     
  Lines        4090     4105      +15     
==========================================
+ Hits         3308     3321      +13     
- Misses        782      784       +2     

☔ View full report in Codecov by Harness.
📢 Have feedback on the report? Share it here.

🚀 New features to boost your workflow:
  • ❄️ Test Analytics: Detect flaky tests, report on failures, and find test suite problems.

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.
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.
@maleadt
maleadt merged commit 5a3557a into main Oct 2, 2026
4 of 5 checks passed
@maleadt
maleadt deleted the tb/cooperative-wait branch October 2, 2026 05:44
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant