Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion Project.toml
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down
11 changes: 8 additions & 3 deletions docs/src/api/streams.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
18 changes: 15 additions & 3 deletions src/hip/HIP.jl
Original file line number Diff line number Diff line change
Expand Up @@ -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")
Expand Down Expand Up @@ -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
Expand Down
35 changes: 15 additions & 20 deletions src/hip/event.jl
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down
92 changes: 23 additions & 69 deletions src/hip/stream.jl
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down
33 changes: 33 additions & 0 deletions test/device/hostcall.jl
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
125 changes: 125 additions & 0 deletions test/hip_core_tests.jl
Original file line number Diff line number Diff line change
Expand Up @@ -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]
Expand Down
Loading