diff --git a/Project.toml b/Project.toml index 0236d89d..27f30e81 100644 --- a/Project.toml +++ b/Project.toml @@ -11,6 +11,7 @@ GPUArrays = "0c68f7d7-f131-5f86-a1c3-88cf8149b2d7" GPUCompiler = "61eb1bfa-7361-4325-ad38-22787b887f55" GPUToolbox = "096a3bc2-3ced-46d0-87f4-dd12716f4bfc" KernelAbstractions = "63c18a36-062a-441e-b654-da1e3ab1ce7c" +KernelInterface = "4ee993da-d684-4d17-a7dd-4e58e78d92bf" LLVM = "929cbde3-209d-540e-8aea-75f648917ca0" LinearAlgebra = "37e2e46d-f89d-539d-b4ee-838fcccc9c8e" OpenCL_jll = "6cb37087-e8b6-5417-8430-1f242f1e46e4" @@ -27,6 +28,7 @@ StaticArrays = "90137ffa-7385-5640-81b9-e52037218182" spirv2clc_jll = "f0274c0c-8c8a-59f1-85b7-f7d60330c5fb" [sources] +KernelInterface = {url = "https://github.com/JuliaGPU/KernelAbstractions.jl", rev = "tb/ki-0.3", subdir = "lib/KernelInterface"} SPIRVIntrinsics = {path = "lib/intrinsics"} [compat] @@ -35,6 +37,7 @@ GPUArrays = "11.2.1" GPUCompiler = "2.9" GPUToolbox = "3.1" KernelAbstractions = "0.9.38" +KernelInterface = "0.3" LLVM = "9.6" LinearAlgebra = "1" OpenCL_jll = "=2024.10.24" @@ -44,7 +47,7 @@ Random = "1" Random123 = "1.7.1" RandomNumbers = "1.6.0" Reexport = "1" -SPIRVIntrinsics = "1.1" +SPIRVIntrinsics = "1.1.3" SPIRV_LLVM_Backend_jll = "23" SPIRV_Tools_jll = "2025.1" StaticArrays = "1" diff --git a/lib/cl/CL.jl b/lib/cl/CL.jl index 2788a72f..6000f8de 100644 --- a/lib/cl/CL.jl +++ b/lib/cl/CL.jl @@ -20,6 +20,7 @@ include("device.jl") include("context.jl") include("cmdqueue.jl") include("event.jl") +include("synchronization.jl") include("memory.jl") include("program.jl") include("kernel.jl") diff --git a/lib/cl/cmdqueue.jl b/lib/cl/cmdqueue.jl index bd82efc5..d3ddeed3 100644 --- a/lib/cl/cmdqueue.jl +++ b/lib/cl/cmdqueue.jl @@ -58,8 +58,10 @@ end # `check_exceptions` surfaces device-side exceptions thrown by kernels that ran on the # queue; disable it where throwing is not an option (e.g. finalizers), in which case the -# exception remains pending until the next check. -function finish(q::CmdQueue; check_exceptions::Bool=true) +# exception remains pending until the next check. `blocking=false` waits cooperatively, +# letting other tasks run, which is not possible in finalizers either. +function finish(q::CmdQueue; check_exceptions::Bool=true, blocking::Bool=true) + blocking || wait_cooperatively(q) OpenCL.check_exceptions(q; rethrow=check_exceptions) return q end diff --git a/lib/cl/synchronization.jl b/lib/cl/synchronization.jl new file mode 100644 index 00000000..88ff3854 --- /dev/null +++ b/lib/cl/synchronization.jl @@ -0,0 +1,106 @@ +# cooperative synchronization +# +# Like CUDA.jl, busy-wait briefly for short operations, then block in the driver on a +# separate thread, so that the waiting task yields to the Julia scheduler, and the thread +# it runs on can do other work and take part in garbage collection. + +using GPUToolbox: @gcsafe_ccall + +# whether the command associated with `evt` has completed. errors count as completion: +# waiting for the event reports them. +function isdone(evt::AbstractEvent) + status = Ref{Cint}() + clGetEventInfo(evt, CL_EVENT_COMMAND_EXECUTION_STATUS, sizeof(Cint), status, C_NULL) + return status[] <= CL_COMPLETE +end + +# before waiting on another thread, which has some overhead, busy-wait for the event, +# initially without even yielding to other tasks. returns whether the event completed. +function spinning_wait(evt::AbstractEvent) + isdone(evt) && return true + for spins in 1:256 + if spins <= 32 + ccall(:jl_cpu_pause, Cvoid, ()) + # allow the GC to run while we're spinning + ccall(:jl_gc_safepoint, Cvoid, ()) + else + yield() + end + isdone(evt) && return true + end + return false +end + +struct SyncRequest + event::AbstractEvent + status::Base.RefValue{cl_int} + done::Base.Event +end + +const MAX_SYNC_THREADS = 4 +const sync_channels = Vector{Channel{SyncRequest}}(undef, MAX_SYNC_THREADS) +const sync_channel_cursor = Threads.Atomic{UInt32}(1) +const sync_channel_lock = ReentrantLock() + +# runs on a thread of its own, waiting for the events it's sent +function synchronization_worker(data::Ptr{Cvoid}) + chan = sync_channels[Int(data)] + while true + req = take!(chan) + GC.@preserve req begin + id = Ref(req.event.id) + req.status[] = @gcsafe_ccall libopencl.clWaitForEvents(1::cl_uint, + id::Ptr{cl_event})::cl_int + end + notify(req.done) + end +end + +@noinline function sync_channel(i::Int) + @lock sync_channel_lock begin + isassigned(sync_channels, i) && return sync_channels[i] + chan = Channel{SyncRequest}(Inf) + sync_channels[i] = chan + + # we don't know the size of uv_thread_t, so reserve enough space + tid = Ref{NTuple{32, UInt8}}(ntuple(_ -> 0x00, 32)) + cb = @cfunction(synchronization_worker, Cvoid, (Ptr{Cvoid},)) + err = @ccall uv_thread_create(tid::Ptr{Cvoid}, cb::Ptr{Cvoid}, Ptr{Cvoid}(i)::Ptr{Cvoid})::Cint + err == 0 || Base.uv_error("uv_thread_create", err) + err = @ccall uv_thread_detach(tid::Ptr{Cvoid})::Cint + err == 0 || Base.uv_error("uv_thread_detach", err) + return chan + end +end + +# wait for `evt` on a worker thread, while the calling task yields +function nonblocking_wait(evt::AbstractEvent) + # sticky per task, so that a task keeps using the same worker, while concurrent tasks + # spread over the workers + i = get!(task_local_storage(), :CLSyncChannel) do + mod1(Int(Threads.atomic_add!(sync_channel_cursor, UInt32(1))), MAX_SYNC_THREADS) + end::Int + chan = isassigned(sync_channels, i) ? sync_channels[i] : sync_channel(i) + + req = SyncRequest(evt, Ref{cl_int}(CL_SUCCESS), Base.Event()) + put!(chan, req) + wait(req.done) + req.status[] == CL_SUCCESS || throw(CLError(req.status[])) + return +end + +""" + cl.wait_cooperatively(q::CmdQueue) + +Wait for all commands queued on `q` to complete, while letting other tasks run. This does +not check for errors or device-side exceptions, `cl.finish` does. +""" +function wait_cooperatively(q::CmdQueue) + evt = Ref{cl_event}() + clEnqueueMarkerWithWaitList(q, 0, C_NULL, evt) + marker = Event(evt[]) + # the marker only completes once the queue has been submitted to the device + clFlush(q) + spinning_wait(marker) || nonblocking_wait(marker) + return +end diff --git a/src/OpenCL.jl b/src/OpenCL.jl index da994572..ca02eb1b 100644 --- a/src/OpenCL.jl +++ b/src/OpenCL.jl @@ -12,6 +12,8 @@ using Preferences import KernelAbstractions: KernelAbstractions +import KernelInterface + using Core: LLVMPtr # library wrappers @@ -49,7 +51,12 @@ include("mapreduce.jl") include("gpuarrays.jl") include("random.jl") -include("OpenCLKernels.jl") +include("OpenCLKernelsOld.jl") import .OpenCLKernels: OpenCLBackend export OpenCLBackend + +# KernelInterface - NOT PUBLIC. Use KernelInterface.get_backend on an CLArray to get the backend +include("OpenCLKernels.jl") +import .OpenCLInterface + end diff --git a/src/OpenCLKernels.jl b/src/OpenCLKernels.jl index b1d4eaa7..99410608 100644 --- a/src/OpenCLKernels.jl +++ b/src/OpenCLKernels.jl @@ -1,9 +1,11 @@ -module OpenCLKernels +module OpenCLInterface using ..OpenCL -using ..OpenCL: @device_override, method_table +using ..OpenCL: @device_override, method_table, kernel_convert, clfunction -import KernelAbstractions as KA +import KernelInterface as KI + +import SPIRVIntrinsics import StaticArrays @@ -12,19 +14,35 @@ import Adapt ## Back-end Definition -export OpenCLBackend +# export OpenCLBackend + -Base.@kwdef struct OpenCLBackend <: KA.GPU +# The platform is part of the backend's configuration. A backend works with the task's +# active device if that is on its platform. Otherwise, work for the backend (allocations, +# copies, compilation and launches) first activates the default device of its platform, +# as `KI.device!` would, so that the arrays it creates can be used afterwards. +Base.@kwdef struct OpenCLBackend <: KI.Backend platform::cl.Platform = cl.platform() end -@noinline function platform_mismatch_warning(expected::cl.Platform, active::cl.Platform) - @warn "OpenCLBackend platform \"$(expected.name)\" is not the active platform \"$(active.name)\"" - return nothing +# the device that `b` works with +function backend_device(b::OpenCLBackend) + cl.platform() == b.platform && return cl.device() + dev = cl.default_device(b.platform) + dev === nothing && throw(ArgumentError("OpenCL platform \"$(b.platform.name)\" has no devices")) + return dev end -function KA.allocate(b::OpenCLBackend, ::Type{T}, dims::Tuple; unified::Bool = false) where T - b.platform === cl.platform() || platform_mismatch_warning(b.platform, cl.platform()) +# make the backend's device the task's active device +@inline function activate(b::OpenCLBackend) + cl.platform() == b.platform || cl.platform!(b.platform) + return +end + +KI.versioninfo(io::IO, ::OpenCLBackend) = OpenCL.versioninfo(io) + +function KI.allocate(b::OpenCLBackend, ::Type{T}, dims::Tuple; unified::Bool = false) where T + activate(b) if unified memory_backend = cl.unified_memory_backend() if memory_backend === cl.USMBackend() @@ -39,40 +57,42 @@ function KA.allocate(b::OpenCLBackend, ::Type{T}, dims::Tuple; unified::Bool = f end end -KA.supports_unified(::OpenCLBackend) = cl.default_memory_backend(cl.device(); unified=true) !== nothing +# OpenCL.jl creates a context per device +context_device(ctx::cl.Context) = ctx == cl.context() ? cl.device() : only(ctx.devices) -KA.get_backend(::CLArray) = OpenCLBackend() -# TODO should be non-blocking -KA.synchronize(::OpenCLBackend) = cl.finish(cl.queue()) -KA.supports_float64(::OpenCLBackend) = in("cl_khr_fp64", cl.device().extensions) - -Adapt.adapt_storage(::OpenCLBackend, a::Array) = Adapt.adapt(CLArray, a) -Adapt.adapt_storage(::OpenCLBackend, a::CLArray) = a -Adapt.adapt_storage(::KA.CPU, a::CLArray) = convert(Array, a) +function KI.get_backend(A::CLArray) + ctx = OpenCL.context(A) + ctx == cl.context() && return OpenCLBackend(cl.platform()) + return OpenCLBackend(context_device(ctx).platform) +end -# `@Const` applies `constify` inside the kernel, where arguments have already been -# converted to device arrays, so the rule has to be registered for `CLDeviceArray` -# rather than for `CLArray`. -Adapt.adapt_storage(::KA.ConstAdaptor, a::CLDeviceArray) = Base.Experimental.Const(a) +function KI.synchronize(b::OpenCLBackend) + activate(b) + cl.finish(cl.queue(); blocking=false) + return +end ## Device Selection # devices are numbered consecutively within the backend's platform, in enumeration order -function KA.ndevices(b::OpenCLBackend) +function KI.ndevices(b::OpenCLBackend) Int(cl.ndevices(b.platform)) end -function KA.device(b::OpenCLBackend) - current = cl.device() - for (i, d) in enumerate(cl.devices(b.platform)) - d == current && return i - end - error("Active OpenCL device $current not found in the OpenCLBackend's platform \"$(b.platform.name)\".") +function device_index(b::OpenCLBackend, dev::cl.Device) + id = findfirst(==(dev), cl.devices(b.platform)) + id === nothing && + throw(ArgumentError("OpenCL device $(dev.name) is not on the backend's platform \"$(b.platform.name)\"")) + return id end -function KA.device!(b::OpenCLBackend, id::Int) - 0 < id <= KA.ndevices(b) || throw(ArgumentError("Device id $id out of bounds.")) +KI.device(b::OpenCLBackend) = device_index(b, backend_device(b)) + +KI.device(b::OpenCLBackend, A::CLArray) = device_index(b, context_device(OpenCL.context(A))) + +function KI.device!(b::OpenCLBackend, id::Int) + 0 < id <= KI.ndevices(b) || throw(ArgumentError("Device id $id out of bounds.")) devs = cl.devices(b.platform) cl.device!(devs[id]) @@ -81,152 +101,178 @@ end ## Memory Operations -function KA.copyto!(::OpenCLBackend, A, B) +function KI.copyto!(b::OpenCLBackend, A, B) + length(A) == length(B) || + throw(ArgumentError("Arrays must have the same length, got $(length(A)) and $(length(B))")) + activate(b) copyto!(A, B) - # TODO: Address device to host copies in jl being synchronizing + return A end +KI.unsafe_free!(A::CLArray) = OpenCL.unsafe_free!(A) -## Kernel Launch - -function KA.mkcontext(kernel::KA.Kernel{OpenCLBackend}, _ndrange, iterspace) - KA.CompilerMetadata{KA.ndrange(kernel), KA.DynamicCheck}(_ndrange, iterspace) -end -function KA.mkcontext(kernel::KA.Kernel{OpenCLBackend}, I, _ndrange, iterspace, - ::Dynamic) where Dynamic - KA.CompilerMetadata{KA.ndrange(kernel), Dynamic}(I, _ndrange, iterspace) -end -function KA.launch_config(kernel::KA.Kernel{OpenCLBackend}, ndrange, workgroupsize) - if ndrange isa Integer - ndrange = (ndrange,) - end - if workgroupsize isa Integer - workgroupsize = (workgroupsize, ) - end +## Kernel Launch - # partition checked that the ndrange's agreed - if KA.ndrange(kernel) <: KA.StaticSize - ndrange = nothing - end - iterspace, dynamic = if KA.workgroupsize(kernel) <: KA.DynamicSize && - workgroupsize === nothing - # use ndrange as preliminary workgroupsize for autotuning - KA.partition(kernel, ndrange, ndrange) - else - KA.partition(kernel, ndrange, workgroupsize) - end +KI.argconvert(::OpenCLBackend, arg) = kernel_convert(arg) - return ndrange, workgroupsize, iterspace, dynamic +function KI.kernel_function(backend::OpenCLBackend, f::F, tt::TT=Tuple{}; name = nothing, kwargs...) where {F,TT} + activate(backend) + # on devices that support it, `clfunction` fixes the sub-group width to + # `cl.sub_group_size(dev)`, as `KI.sub_group_size` promises + kern = clfunction(f, tt; name, kwargs...) + KI.Kernel{OpenCLBackend, typeof(kern)}(backend, kern) end -function threads_to_workgroupsize(threads, ndrange) - total = 1 - return map(ndrange) do n - x = min(div(threads, total), n) - total *= x - return x - end +# the context that a kernel was compiled for +function kernel_context(kernel::KI.Kernel{OpenCLBackend}) + ctx = Ref{cl.cl_context}() + cl.clGetKernelInfo(kernel.kern.fun, cl.CL_KERNEL_CONTEXT, sizeof(cl.cl_context), ctx, C_NULL) + return ctx[] end -function (obj::KA.Kernel{OpenCLBackend})(args...; ndrange=nothing, workgroupsize=nothing) - obj.backend.platform === cl.platform() || platform_mismatch_warning(obj.backend.platform, cl.platform()) +function kernel_device(kernel::KI.Kernel{OpenCLBackend}) + ctx = kernel_context(kernel) + ctx == cl.context().id && return cl.device() + return context_device(cl.Context(ctx; retain=true)) +end - ndrange, workgroupsize, iterspace, dynamic = - KA.launch_config(obj, ndrange, workgroupsize) +@noinline function throw_device_mismatch(kernel) + throw(ArgumentError("Cannot launch a kernel compiled for $(kernel_device(kernel).name) on $(cl.device().name)")) +end - # this might not be the final context, since we may tune the workgroupsize - ctx = KA.mkcontext(obj, ndrange, iterspace) - kernel = @opencl launch=false obj.f(ctx, args...) +function KI.launch(kernel::KI.Kernel{OpenCLBackend}, groups::Dims{3}, items::Dims{3}, args::Vararg{Any, N}; kwargs...) where {N} + activate(kernel.backend) + kernel_context(kernel) == cl.context().id || throw_device_mismatch(kernel) + kernel.kern(args...; local_size = items, global_size = items .* groups, kwargs...) + return +end - # figure out the optimal workgroupsize automatically - if KA.workgroupsize(obj) <: KA.DynamicSize && workgroupsize === nothing - wg_info = cl.work_group_info(kernel.fun, cl.device()) - wg_size_nd = threads_to_workgroupsize(wg_info.size, ndrange) - iterspace, dynamic = KA.partition(obj, ndrange, wg_size_nd) - ctx = KA.mkcontext(obj, ndrange, iterspace) - end +function KI.max_work_group_size(kernel::KI.Kernel{OpenCLBackend})::Int + wginfo = cl.work_group_info(kernel.kern.fun, kernel_device(kernel)) + Int(wginfo.size) +end - groups = length(KA.blocks(iterspace)) - items = length(KA.workitems(iterspace)) - if groups == 0 - return nothing +## Device Properties + +# querying the device allocates, so cache what launches and kernels need. the cache is +# keyed on the device, because the task-local device can be switched. +const DeviceProperties = @NamedTuple{ + max_work_group_size::Int, max_work_group_dims::NTuple{3, Int}, compute_units::Int, + float64::Bool, float16::Bool, unified::Bool, + # 0 if the device doesn't support sub-groups of a fixed width + sub_group_size::Int, sub_group_shuffle::Bool, +} +function device_properties(dev::cl.Device) + cache = get!(task_local_storage(), :CLDeviceProperties) do + Dict{cl.Device, DeviceProperties}() + end::Dict{cl.Device, DeviceProperties} + return get!(cache, dev) do + sizes = dev.max_work_item_size + extensions = dev.extensions + # the sub-group width is only fixed for kernels that request it, which `clfunction` + # does for devices with `cl_intel_required_subgroup_size` + fixed_sub_groups = cl.sub_groups_supported(dev) && + "cl_intel_required_subgroup_size" in extensions + (; max_work_group_size = Int(dev.max_work_group_size), + max_work_group_dims = ntuple(d -> d <= length(sizes) ? Int(sizes[d]) : 1, 3), + compute_units = Int(dev.max_compute_units), + float64 = "cl_khr_fp64" in extensions, + float16 = "cl_khr_fp16" in extensions, + unified = cl.default_memory_backend(dev; unified=true) !== nothing, + sub_group_size = fixed_sub_groups ? cl.sub_group_size(dev) : 0, + sub_group_shuffle = fixed_sub_groups && "cl_khr_subgroup_shuffle" in extensions) end +end - # Launch kernel - global_size = groups * items - local_size = items - kernel(ctx, args...; global_size, local_size) +device_properties(b::OpenCLBackend) = device_properties(backend_device(b)) - return nothing +KI.max_work_group_size(b::OpenCLBackend)::Int = device_properties(b).max_work_group_size +KI.max_work_group_dims(b::OpenCLBackend)::NTuple{3, Int} = device_properties(b).max_work_group_dims +# OpenCL doesn't limit the number of work-groups, only the global size (to `size_t`) +function KI.max_num_groups(b::OpenCLBackend)::NTuple{3, Int} + return typemax(Int) .รท KI.max_work_group_dims(b) +end +KI.multiprocessor_count(b::OpenCLBackend)::Int = device_properties(b).compute_units + +KI.supports_float64(b::OpenCLBackend) = device_properties(b).float64 +KI.supports_unified(b::OpenCLBackend) = device_properties(b).unified +# 32-bit integer atomics are core OpenCL; float atomics fall back to compare-and-swap +KI.supports_atomics(::OpenCLBackend) = true + +KI.supports_subgroups(b::OpenCLBackend) = device_properties(b).sub_group_size > 0 +KI.sub_group_size(b::OpenCLBackend)::Int = device_properties(b).sub_group_size +function KI.supports_shuffle(b::OpenCLBackend, ::Type{T}) where {T} + props = device_properties(b) + props.sub_group_shuffle || return false + T in SPIRVIntrinsics.gentypes || return false + T === Float64 && return props.float64 + T === Float16 && return props.float16 + return true end - ## Indexing Functions +## COV_EXCL_START -@device_override @inline function KA.__index_Local_Linear(ctx) - return get_local_id(1) -end +# computed with `% T`, which unlike `T(x)` has no error path. KernelInterface derives the +# global queries from these. -@device_override @inline function KA.__index_Group_Linear(ctx) - return get_group_id(1) +@device_override @inline function KI.get_local_id(::Type{T}) where {T} + return (; x = get_local_id(1) % T, y = get_local_id(2) % T, z = get_local_id(3) % T) end -@device_override @inline function KA.__index_Global_Linear(ctx) - #return get_global_id(1) # JuliaGPU/OpenCL.jl#346 - I = KA.__index_Global_Cartesian(ctx) - @inbounds LinearIndices(KA.__ndrange(ctx))[I] +@device_override @inline function KI.get_group_id(::Type{T}) where {T} + return (; x = get_group_id(1) % T, y = get_group_id(2) % T, z = get_group_id(3) % T) end -@device_override @inline function KA.__index_Local_Cartesian(ctx) - @inbounds KA.workitems(KA.__iterspace(ctx))[get_local_id(1)] +@device_override @inline function KI.get_local_size(::Type{T}) where {T} + return (; x = get_local_size(1) % T, y = get_local_size(2) % T, z = get_local_size(3) % T) end -@device_override @inline function KA.__index_Group_Cartesian(ctx) - @inbounds KA.blocks(KA.__iterspace(ctx))[get_group_id(1)] +@device_override @inline function KI.get_num_groups(::Type{T}) where {T} + return (; x = get_num_groups(1) % T, y = get_num_groups(2) % T, z = get_num_groups(3) % T) end -@device_override @inline function KA.__index_Global_Cartesian(ctx) - return @inbounds KA.expand(KA.__iterspace(ctx), get_group_id(1), get_local_id(1)) -end +# OpenCL's sub-group queries already have KernelInterface's semantics: the last sub-group +# of a work-group can be partial, and `get_sub_group_size` counts the work-items present -@device_override @inline function KA.__validindex(ctx) - if KA.__dynamic_checkbounds(ctx) - I = KA.__index_Global_Cartesian(ctx) - return I in KA.__ndrange(ctx) - else - return true - end -end +@device_override KI.get_sub_group_size(::Type{T}) where {T} = get_sub_group_size() % T + +@device_override KI.get_max_sub_group_size(::Type{T}) where {T} = get_max_sub_group_size() % T + +@device_override KI.get_num_sub_groups(::Type{T}) where {T} = get_num_sub_groups() % T +@device_override KI.get_sub_group_id(::Type{T}) where {T} = get_sub_group_id() % T + +@device_override KI.get_sub_group_local_id(::Type{T}) where {T} = get_sub_group_local_id() % T ## Shared and Scratch Memory -@device_override @inline function KA.SharedMemory(::Type{T}, ::Val{Dims}, ::Val{Id}) where {T, Dims, Id} +@device_override @inline function KI.localmemory(::Type{T}, ::Val{Dims}) where {T, Dims} ptr = OpenCL.emit_localmemory(T, Val(prod(Dims))) CLDeviceArray(Dims, ptr) end -@device_override @inline function KA.Scratchpad(ctx, ::Type{T}, ::Val{Dims}) where {T, Dims} - StaticArrays.MArray{KA.__size(Dims), T}(undef) -end - - ## Synchronization and Printing -@device_override @inline function KA.__synchronize() +@device_override @inline function KI.barrier() work_group_barrier(OpenCL.LOCAL_MEM_FENCE | OpenCL.GLOBAL_MEM_FENCE) end -@device_override @inline function KA.__print(args...) - OpenCL._print(args...) +@device_override @inline function KI.sub_group_barrier() + sub_group_barrier(OpenCL.LOCAL_MEM_FENCE | OpenCL.GLOBAL_MEM_FENCE) end +# out-of-range source lanes give an undefined value, as KernelInterface allows +@device_override function KI.shfl_down(val::T, offset::Integer) where T + sub_group_shuffle(val, get_sub_group_local_id() % UInt32 + offset % UInt32) +end -## Other - -KA.argconvert(::KA.Kernel{OpenCLBackend}, arg) = OpenCL.kernel_convert(arg) +@device_override @inline function KI._print(args...) + OpenCL._print(args...) +end +## COV_EXCL_STOP end diff --git a/src/OpenCLKernelsOld.jl b/src/OpenCLKernelsOld.jl new file mode 100644 index 00000000..b1d4eaa7 --- /dev/null +++ b/src/OpenCLKernelsOld.jl @@ -0,0 +1,232 @@ +module OpenCLKernels + +using ..OpenCL +using ..OpenCL: @device_override, method_table + +import KernelAbstractions as KA + +import StaticArrays + +import Adapt + + +## Back-end Definition + +export OpenCLBackend + +Base.@kwdef struct OpenCLBackend <: KA.GPU + platform::cl.Platform = cl.platform() +end + +@noinline function platform_mismatch_warning(expected::cl.Platform, active::cl.Platform) + @warn "OpenCLBackend platform \"$(expected.name)\" is not the active platform \"$(active.name)\"" + return nothing +end + +function KA.allocate(b::OpenCLBackend, ::Type{T}, dims::Tuple; unified::Bool = false) where T + b.platform === cl.platform() || platform_mismatch_warning(b.platform, cl.platform()) + if unified + memory_backend = cl.unified_memory_backend() + if memory_backend === cl.USMBackend() + return CLArray{T, length(dims), cl.UnifiedSharedMemory}(undef, dims) + elseif memory_backend === cl.SVMBackend() + return CLArray{T, length(dims), cl.SharedVirtualMemory}(undef, dims) + else + throw(ArgumentError("Unified memory not supported")) + end + else + return CLArray{T}(undef, dims) + end +end + +KA.supports_unified(::OpenCLBackend) = cl.default_memory_backend(cl.device(); unified=true) !== nothing + +KA.get_backend(::CLArray) = OpenCLBackend() +# TODO should be non-blocking +KA.synchronize(::OpenCLBackend) = cl.finish(cl.queue()) +KA.supports_float64(::OpenCLBackend) = in("cl_khr_fp64", cl.device().extensions) + +Adapt.adapt_storage(::OpenCLBackend, a::Array) = Adapt.adapt(CLArray, a) +Adapt.adapt_storage(::OpenCLBackend, a::CLArray) = a +Adapt.adapt_storage(::KA.CPU, a::CLArray) = convert(Array, a) + +# `@Const` applies `constify` inside the kernel, where arguments have already been +# converted to device arrays, so the rule has to be registered for `CLDeviceArray` +# rather than for `CLArray`. +Adapt.adapt_storage(::KA.ConstAdaptor, a::CLDeviceArray) = Base.Experimental.Const(a) + +## Device Selection + +# devices are numbered consecutively within the backend's platform, in enumeration order + +function KA.ndevices(b::OpenCLBackend) + Int(cl.ndevices(b.platform)) +end + +function KA.device(b::OpenCLBackend) + current = cl.device() + for (i, d) in enumerate(cl.devices(b.platform)) + d == current && return i + end + error("Active OpenCL device $current not found in the OpenCLBackend's platform \"$(b.platform.name)\".") +end + +function KA.device!(b::OpenCLBackend, id::Int) + 0 < id <= KA.ndevices(b) || throw(ArgumentError("Device id $id out of bounds.")) + devs = cl.devices(b.platform) + + cl.device!(devs[id]) + return nothing +end + +## Memory Operations + +function KA.copyto!(::OpenCLBackend, A, B) + copyto!(A, B) + # TODO: Address device to host copies in jl being synchronizing +end + + +## Kernel Launch + +function KA.mkcontext(kernel::KA.Kernel{OpenCLBackend}, _ndrange, iterspace) + KA.CompilerMetadata{KA.ndrange(kernel), KA.DynamicCheck}(_ndrange, iterspace) +end +function KA.mkcontext(kernel::KA.Kernel{OpenCLBackend}, I, _ndrange, iterspace, + ::Dynamic) where Dynamic + KA.CompilerMetadata{KA.ndrange(kernel), Dynamic}(I, _ndrange, iterspace) +end + +function KA.launch_config(kernel::KA.Kernel{OpenCLBackend}, ndrange, workgroupsize) + if ndrange isa Integer + ndrange = (ndrange,) + end + if workgroupsize isa Integer + workgroupsize = (workgroupsize, ) + end + + # partition checked that the ndrange's agreed + if KA.ndrange(kernel) <: KA.StaticSize + ndrange = nothing + end + + iterspace, dynamic = if KA.workgroupsize(kernel) <: KA.DynamicSize && + workgroupsize === nothing + # use ndrange as preliminary workgroupsize for autotuning + KA.partition(kernel, ndrange, ndrange) + else + KA.partition(kernel, ndrange, workgroupsize) + end + + return ndrange, workgroupsize, iterspace, dynamic +end + +function threads_to_workgroupsize(threads, ndrange) + total = 1 + return map(ndrange) do n + x = min(div(threads, total), n) + total *= x + return x + end +end + +function (obj::KA.Kernel{OpenCLBackend})(args...; ndrange=nothing, workgroupsize=nothing) + obj.backend.platform === cl.platform() || platform_mismatch_warning(obj.backend.platform, cl.platform()) + + ndrange, workgroupsize, iterspace, dynamic = + KA.launch_config(obj, ndrange, workgroupsize) + + # this might not be the final context, since we may tune the workgroupsize + ctx = KA.mkcontext(obj, ndrange, iterspace) + kernel = @opencl launch=false obj.f(ctx, args...) + + # figure out the optimal workgroupsize automatically + if KA.workgroupsize(obj) <: KA.DynamicSize && workgroupsize === nothing + wg_info = cl.work_group_info(kernel.fun, cl.device()) + wg_size_nd = threads_to_workgroupsize(wg_info.size, ndrange) + iterspace, dynamic = KA.partition(obj, ndrange, wg_size_nd) + ctx = KA.mkcontext(obj, ndrange, iterspace) + end + + groups = length(KA.blocks(iterspace)) + items = length(KA.workitems(iterspace)) + + if groups == 0 + return nothing + end + + # Launch kernel + global_size = groups * items + local_size = items + kernel(ctx, args...; global_size, local_size) + + return nothing +end + + +## Indexing Functions + +@device_override @inline function KA.__index_Local_Linear(ctx) + return get_local_id(1) +end + +@device_override @inline function KA.__index_Group_Linear(ctx) + return get_group_id(1) +end + +@device_override @inline function KA.__index_Global_Linear(ctx) + #return get_global_id(1) # JuliaGPU/OpenCL.jl#346 + I = KA.__index_Global_Cartesian(ctx) + @inbounds LinearIndices(KA.__ndrange(ctx))[I] +end + +@device_override @inline function KA.__index_Local_Cartesian(ctx) + @inbounds KA.workitems(KA.__iterspace(ctx))[get_local_id(1)] +end + +@device_override @inline function KA.__index_Group_Cartesian(ctx) + @inbounds KA.blocks(KA.__iterspace(ctx))[get_group_id(1)] +end + +@device_override @inline function KA.__index_Global_Cartesian(ctx) + return @inbounds KA.expand(KA.__iterspace(ctx), get_group_id(1), get_local_id(1)) +end + +@device_override @inline function KA.__validindex(ctx) + if KA.__dynamic_checkbounds(ctx) + I = KA.__index_Global_Cartesian(ctx) + return I in KA.__ndrange(ctx) + else + return true + end +end + + +## Shared and Scratch Memory + +@device_override @inline function KA.SharedMemory(::Type{T}, ::Val{Dims}, ::Val{Id}) where {T, Dims, Id} + ptr = OpenCL.emit_localmemory(T, Val(prod(Dims))) + CLDeviceArray(Dims, ptr) +end + +@device_override @inline function KA.Scratchpad(ctx, ::Type{T}, ::Val{Dims}) where {T, Dims} + StaticArrays.MArray{KA.__size(Dims), T}(undef) +end + + +## Synchronization and Printing + +@device_override @inline function KA.__synchronize() + work_group_barrier(OpenCL.LOCAL_MEM_FENCE | OpenCL.GLOBAL_MEM_FENCE) +end + +@device_override @inline function KA.__print(args...) + OpenCL._print(args...) +end + + +## Other + +KA.argconvert(::KA.Kernel{OpenCLBackend}, arg) = OpenCL.kernel_convert(arg) + +end diff --git a/src/util.jl b/src/util.jl index 725fc1f4..8fbe58d8 100644 --- a/src/util.jl +++ b/src/util.jl @@ -73,7 +73,7 @@ function versioninfo(io::IO=stdout) end for pkg in [:GPUArrays, :GPUCompiler, ("63c18a36-062a-441e-b654-da1e3ab1ce7c", "KernelAbstractions"), - :LLVM, :SPIRVIntrinsics, ("627d6b7a-bbe6-5189-83e7-98cc0a5aeadd", "pocl_jll"), + :KernelInterface, :LLVM, :SPIRVIntrinsics, ("627d6b7a-bbe6-5189-83e7-98cc0a5aeadd", "pocl_jll"), ("59abdad9-3cfc-5436-8271-411e8cad6b82", "pocl_next_jll")] name, mod = get_module(pkg) isnothing(mod) && continue diff --git a/test/Project.toml b/test/Project.toml index 14eef6fd..6d154cb2 100644 --- a/test/Project.toml +++ b/test/Project.toml @@ -9,6 +9,7 @@ IOCapture = "b5f81e59-6552-4d32-b1f0-c071b021bf89" InteractiveUtils = "b77e0a4c-d291-57a0-90e8-8db25a27a240" JLD2 = "033835bb-8acc-5ee8-8aae-3f567f8a3819" KernelAbstractions = "63c18a36-062a-441e-b654-da1e3ab1ce7c" +KernelInterface = "4ee993da-d684-4d17-a7dd-4e58e78d92bf" LinearAlgebra = "37e2e46d-f89d-539d-b4ee-838fcccc9c8e" OpenCL = "08131aa3-fb12-5dee-8b74-c09406e224a2" ParallelTestRunner = "d3525ed8-44d0-4b2c-a655-542cee43accc" @@ -29,6 +30,7 @@ pocl_jll = "627d6b7a-bbe6-5189-83e7-98cc0a5aeadd" pocl_next_jll = "59abdad9-3cfc-5436-8271-411e8cad6b82" [sources] +KernelInterface = {url = "https://github.com/JuliaGPU/KernelAbstractions.jl", rev = "tb/ki-0.3", subdir = "lib/KernelInterface"} OpenCL = {path = ".."} SPIRVIntrinsics = {path = "../lib/intrinsics"} diff --git a/test/kernelinterface.jl b/test/kernelinterface.jl new file mode 100644 index 00000000..63511d0b --- /dev/null +++ b/test/kernelinterface.jl @@ -0,0 +1,83 @@ +import KernelInterface +import KernelInterface as KI +using OpenCL.OpenCLInterface + +include(joinpath(dirname(pathof(KernelInterface)), "..", "test", "testsuite.jl")) + +Testsuite.testsuite(OpenCLInterface.OpenCLBackend(), CLArray) + +function ki_fill_kernel(a, val) + i = KI.get_global_id().x + if i <= length(a) + @inbounds a[i] = val + end + return +end + +@testset "backend platform" begin + backend = OpenCLInterface.OpenCLBackend() + a = KI.zeros(backend, Int32, 4) + @test KI.get_backend(a) == backend + + # a kernel belongs to the device it was compiled for + kernel = KI.@launch backend launch=false ki_fill_kernel(a, Int32(1)) + cl.context!(cl.Context(cl.device())) do + @test_throws ArgumentError kernel(a, Int32(1); ndrange = 4) + end + kernel(a, Int32(1); ndrange = 4) + @test Array(a) == ones(Int32, 4) + + # a backend for another platform activates it + others = filter(!=(cl.platform()), cl.platforms()) + if !isempty(others) + platform, device = cl.platform(), cl.device() + try + other = OpenCLInterface.OpenCLBackend(; platform = first(others)) + @test KI.device(other) == 1 + b = KI.zeros(other, Int32, 4) + @test cl.platform() == first(others) + @test KI.get_backend(b) == other + @test KI.device(other, b) == KI.device(other) + KI.@launch other ndrange = 4 ki_fill_kernel(b, Int32(2)) + @test Array(b) == fill(Int32(2), 4) + + # arrays keep their platform + cl.device!(device) + @test KI.get_backend(b) == other + @test KI.get_backend(a) == backend + finally + cl.platform!(platform) + cl.device!(device) + end + end +end + +function ki_slow_kernel(a, iters) + acc = UInt32(KI.get_global_id().x) + for k in UInt32(1):iters + acc = acc * 0x0019660d + k + end + @inbounds a[1] = acc + return +end + +@testset "cooperative synchronize" begin + backend = OpenCLInterface.OpenCLBackend() + a = KI.zeros(backend, UInt32, 1) + KI.@launch backend ki_slow_kernel(a, UInt32(1)) + KI.synchronize(backend) + + # another task on this thread gets to run while `synchronize` waits for the device + done = Ref(false) + ticks = Ref(0) + task = @async while !done[] + ticks[] += 1 + yield() + end + KI.@launch backend ki_slow_kernel(a, UInt32(2)^24) + KI.synchronize(backend) + during = ticks[] + done[] = true + wait(task) + @test during > 0 +end