diff --git a/Project.toml b/Project.toml index 8bc9180ef..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" +GPUToolbox = "3.3.1" 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..9b8ea4a39 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 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 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..87c417b7d 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") @@ -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/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/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 8e38241a3..63d236ce9 100644 --- a/test/hip_core_tests.jl +++ b/test/hip_core_tests.jl @@ -22,6 +22,131 @@ Random.seed!(1) @test t >= 0 end +@testset "cooperative synchronization" begin + # 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 + 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 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 + wait(t) + end + end + + # 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 = streams[1] + @test gated(s) do + open_gate_during(() -> AMDGPU.synchronize(s)) && HIP.isdone(s) + end == (true, false) + end + + 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 = streams[3] + @test gated(s) do + HIP.record(event) + open_gate_during(() -> HIP.synchronize(event)) && HIP.isdone(event) + end == (true, false) + end + + 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 gated(s) do + open_gate_during(() -> AMDGPU.synchronize(s)) && HIP.isdone(s) + end == (true, false) + end + + # 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!(() -> (@roc stream=s noop_kernel()), other) + HIP.synchronize(s; spin=false) + @test HIP.isdone(s) + AMDGPU.device!(HIP.device_synchronize, other) + @test AMDGPU.device() == dev + end +end + if length(AMDGPU.devices()) > 1 @testset "HIP Peer Access" begin dev1, dev2 = AMDGPU.devices()[1:2]