Skip to content
Merged
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
@@ -1,6 +1,6 @@
name = "GPUToolbox"
uuid = "096a3bc2-3ced-46d0-87f4-dd12716f4bfc"
version = "3.3.1"
version = "3.3.2"

[deps]
LLVM = "929cbde3-209d-540e-8aea-75f648917ca0"
Expand Down
45 changes: 41 additions & 4 deletions src/synchronization.jl
Original file line number Diff line number Diff line change
Expand Up @@ -195,13 +195,50 @@ end
# loop to be available). the calling thread is adopted by Julia, which is safe as long as
# driver calls that may be waiting on the callback are GC-safe.

# a spin lock-based condition, so that signalling from a driver thread never blocks it in
# Julia's scheduler (as a `ReentrantLock` could)
# a spin lock whose waiters give up their CPU when they have to wait for long. the lock of a
# completion is held by the driver's thread while it wakes up the waiting task, which
# immediately takes the lock again. if that task runs on the same CPU, it preempts the
# driver's thread, and would spin until the OS switches back. when all cores are busy (e.g.,
# with PoCL executing the next kernel), that can take hundreds of µs.
struct YieldingSpinLock <: Base.AbstractLock
lock::Threads.SpinLock
end
YieldingSpinLock() = YieldingSpinLock(Threads.SpinLock())

const LOCK_SPIN_ITERATIONS = 100

function Base.lock(l::YieldingSpinLock)
i = 0
while !trylock(l.lock)
if i < LOCK_SPIN_ITERATIONS
ccall(:jl_cpu_pause, Cvoid, ())
i += 1
else
yield_cpu()
end
GC.safepoint()
end
return
end
Base.trylock(l::YieldingSpinLock) = trylock(l.lock)
Base.unlock(l::YieldingSpinLock) = unlock(l.lock)
Base.islocked(l::YieldingSpinLock) = islocked(l.lock)
Base.assert_havelock(l::YieldingSpinLock) = Base.assert_havelock(l.lock)

# let the OS run another thread on this CPU
@static if Sys.iswindows()
yield_cpu() = ccall((:SwitchToThread, "kernel32"), stdcall, Cint, ())
else
yield_cpu() = ccall(:sched_yield, Cint, ())
end

# a condition protected by a spin lock, so that signalling from a driver thread never blocks
# it in Julia's scheduler (as a `ReentrantLock` could)
mutable struct Completion
const cond::Base.ThreadSynchronizer
const cond::Base.GenericCondition{YieldingSpinLock}
@atomic state::Int # one of the constants below

Completion() = new(Base.ThreadSynchronizer(), PENDING)
Completion() = new(Base.GenericCondition(YieldingSpinLock()), PENDING)
end
const PENDING = 0
const SIGNALLED = 1
Expand Down
64 changes: 64 additions & 0 deletions test/synchronization.jl
Original file line number Diff line number Diff line change
Expand Up @@ -333,3 +333,67 @@ end
GC.gc(); GC.gc()
@test isassigned(ret) && ret[] == Some(:waited)
end

# run `f` on a thread that is not managed by Julia
mutable struct ForeignCall
const f::Any
@atomic done::Bool
@atomic failed::Bool
end
const foreign_calls = ForeignCall[] # keep them alive
function foreign_call(arg::Ptr{Cvoid})
call = unsafe_pointer_to_objref(arg)::ForeignCall
try
call.f()
catch
@atomic call.failed = true
end
@atomic call.done = true
return
end
function on_foreign_thread(f)
call = ForeignCall(f, false, false)
push!(foreign_calls, call)
tid = Ref{NTuple{32, UInt8}}(ntuple(_ -> 0x0, 32))
cb = @cfunction(foreign_call, Cvoid, (Ptr{Cvoid},))
err = ccall(:uv_thread_create, Cint, (Ptr{Cvoid}, Ptr{Cvoid}, Ptr{Cvoid}),
tid, cb, pointer_from_objref(call))
@test err == 0
ccall(:uv_thread_detach, Cint, (Ptr{Cvoid},), tid)
return call
end

@testset "completion lock" begin
# a task that is woken up while the notifying thread still holds the lock waits for it
# to be released
cond = Base.GenericCondition(GPUToolbox.YieldingSpinLock())
signalled = Ref(false)
waiter = @async @lock cond begin
while !signalled[]
wait(cond)
end
islocked(cond)
end
call = on_foreign_thread() do
@lock cond begin
signalled[] = true
notify(cond)
@gcsafe_ccall uv_sleep(10::Cuint)::Cvoid
end
end
@test timedwait(() -> istaskdone(waiter), 30) === :ok
@test fetch(waiter)
@test timedwait(() -> @atomic(call.done), 30) === :ok
@test !@atomic(call.failed)

# a driver thread waiting for the lock lets the garbage collector run
c = GPUToolbox.Completion()
GC.@preserve c begin
lock(c.cond)
signal_later(pointer_from_objref(c), 0)
@gcsafe_ccall uv_sleep(10::Cuint)::Cvoid
GC.gc()
unlock(c.cond)
@test timedwait(() -> (@atomic c.state) == GPUToolbox.RELEASED, 30) === :ok
end
end
Loading