diff --git a/Project.toml b/Project.toml index 815e411..0f8e694 100644 --- a/Project.toml +++ b/Project.toml @@ -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" diff --git a/src/synchronization.jl b/src/synchronization.jl index 62e98b0..38c558b 100644 --- a/src/synchronization.jl +++ b/src/synchronization.jl @@ -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 diff --git a/test/synchronization.jl b/test/synchronization.jl index 1720f10..2e1dd8f 100644 --- a/test/synchronization.jl +++ b/test/synchronization.jl @@ -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