diff --git a/.buildkite/pipeline.yml b/.buildkite/pipeline.yml index 7cd9f046..657bc6b2 100644 --- a/.buildkite/pipeline.yml +++ b/.buildkite/pipeline.yml @@ -24,6 +24,11 @@ steps: agents: queue: "oneapi" commands: | + if [[ "{{matrix.julia}}" == "1.10" ]]; then + # XXX: Julia 1.10 ignores [sources]; develop the unregistered KernelInterface instead + git clone --depth 1 --branch tb/ki-0.3 https://github.com/JuliaGPU/KernelAbstractions.jl ka + julia --project -e 'using Pkg; Pkg.develop(path="ka/lib/KernelInterface")' + fi julia --project=deps deps/build_ci.jl if: | build.message !~ /\[skip [^\]]*(tests|julia)/ && diff --git a/.github/workflows/docs.yml b/.github/workflows/docs.yml index c853602f..48a2367b 100644 --- a/.github/workflows/docs.yml +++ b/.github/workflows/docs.yml @@ -31,5 +31,11 @@ jobs: with: version: 'lts' - uses: julia-actions/cache@v3 + # XXX: Julia 1.10 ignores [sources]; develop the unregistered KernelInterface instead + - run: | + git clone --depth 1 --branch tb/ki-0.3 https://github.com/JuliaGPU/KernelAbstractions.jl ka + for project in . docs; do + julia --project=$project -e 'using Pkg; Pkg.develop(path="ka/lib/KernelInterface")' + done - uses: julia-actions/julia-buildpkg@latest - run: julia --project=docs/ docs/make.jl diff --git a/.gitlab-ci.yml b/.gitlab-ci.yml index cf54dd92..6aa57301 100644 --- a/.gitlab-ci.yml +++ b/.gitlab-ci.yml @@ -90,6 +90,15 @@ stages: [build, test, coverage] script: - !reference [.aurora-env, script] - julia --color=yes --project=deps deps/build_local.jl + # XXX: Julia 1.10 ignores [sources]; develop the unregistered KernelInterface instead + # (the checkout travels to the test job as an artifact, like test/Manifest.toml) + - | + if [ "$JULIA_VERSION" = "1.10" ]; then + rm -rf ka + git clone --depth 1 --branch tb/ki-0.3 https://github.com/JuliaGPU/KernelAbstractions.jl ka + julia --color=yes --project=. -e 'using Pkg; Pkg.develop(path="ka/lib/KernelInterface")' + julia --color=yes --project=test -e 'using Pkg; Pkg.develop([PackageSpec(path="."), PackageSpec(path="ka/lib/KernelInterface")])' + fi # Instantiate (and thereby precompile) both environments here so the 1 h batch job # spends its walltime on tests, not on Pkg. Manifests are gitignored, so the test # env must be resolved here with oneAPI dev'ed at the checkout (path "..") — a plain @@ -113,6 +122,7 @@ stages: [build, test, coverage] - LocalPreferences.toml - test/LocalPreferences.toml - test/Manifest.toml + - ka/lib/KernelInterface expire_in: 1 week .test:lts: diff --git a/Project.toml b/Project.toml index 88b57ba4..3c008c85 100644 --- a/Project.toml +++ b/Project.toml @@ -13,6 +13,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" Libdl = "8f399da3-3557-5675-b5ff-fb832c97cbdb" LinearAlgebra = "37e2e46d-f89d-539d-b4ee-838fcccc9c8e" @@ -32,6 +33,9 @@ oneAPI_Level_Zero_Headers_jll = "f4bc562b-d309-54f8-9efb-476e56f0410d" oneAPI_Level_Zero_Loader_jll = "13eca655-d68d-5b81-8367-6d99d727ab01" oneAPI_Support_jll = "b049733a-a71d-5ed3-8eba-7d323ac00b36" +[sources] +KernelInterface = {url = "https://github.com/JuliaGPU/KernelAbstractions.jl", rev = "tb/ki-0.3", subdir = "lib/KernelInterface"} + [compat] AbstractFFTs = "1.5.0" AcceleratedKernels = "0.3.1, 0.4" @@ -42,11 +46,12 @@ GPUArrays = "11.5.14" GPUCompiler = "2.9" GPUToolbox = "3.1" KernelAbstractions = "0.9.39" +KernelInterface = "0.3" LLVM = "6, 7, 8, 9" NEO_jll = "=26.18.38308" PrecompileTools = "1" Preferences = "1" -SPIRVIntrinsics = "1" +SPIRVIntrinsics = "1.1.3" SPIRV_LLVM_Backend_jll = "23" SPIRV_LLVM_Translator_jll = "23" SPIRV_Tools_jll = "2025.4.0" diff --git a/lib/level-zero/module.jl b/lib/level-zero/module.jl index 91cc0dfc..6e0635df 100644 --- a/lib/level-zero/module.jl +++ b/lib/level-zero/module.jl @@ -85,13 +85,17 @@ mutable struct ZeKernel # Read on every launch by the scratch hedge, so it must not cost an API call. spill::Int + # cached maxGroupSize, seeded by `properties`; -1 while unqueried, and 0 without the + # MAX_GROUP_SIZE extension. Read on every KernelInterface launch. + max_group_size::Int + function ZeKernel(mod, name) GC.@preserve name begin desc_ref = Ref(ze_kernel_desc_t(; pKernelName=pointer(name))) handle_ref = Ref{ze_kernel_handle_t}() zeKernelCreate(mod, desc_ref, handle_ref) end - obj = new(mod, handle_ref[], ReentrantLock(), -1) + obj = new(mod, handle_ref[], ReentrantLock(), -1, -1) finalizer(obj) do obj zeKernelDestroy(obj) @@ -268,6 +272,8 @@ function properties(kernel::ZeKernel) props = props_ref[] kernel.spill = Int(props.spillMemSize) + kernel.max_group_size = max_group_size_props_ref === nothing ? 0 : + Int(max_group_size_props_ref[].maxGroupSize) return ( numKernelArgs=Int(props.numKernelArgs), requiredGroupSize=ZeDim3(props.requiredGroupSizeX, @@ -295,6 +301,13 @@ function spill_mem_size(kernel::ZeKernel) return s >= 0 ? s : Int(properties(kernel).spillMemSize) end +# Cached access to a kernel's maxGroupSize, `missing` without the MAX_GROUP_SIZE extension. +function max_group_size(kernel::ZeKernel) + s = kernel.max_group_size + s < 0 && (s = coalesce(properties(kernel).maxGroupSize, 0)) + return s > 0 ? s : missing +end + ## execution diff --git a/lib/level-zero/oneL0.jl b/lib/level-zero/oneL0.jl index 82e33eec..c36b9eee 100644 --- a/lib/level-zero/oneL0.jl +++ b/lib/level-zero/oneL0.jl @@ -124,6 +124,7 @@ end include("context.jl") include("cmdqueue.jl") include("cmdlist.jl") +include("synchronization.jl") include("fence.jl") include("event.jl") include("barrier.jl") diff --git a/lib/level-zero/synchronization.jl b/lib/level-zero/synchronization.jl new file mode 100644 index 00000000..809a03b8 --- /dev/null +++ b/lib/level-zero/synchronization.jl @@ -0,0 +1,195 @@ +# cooperative synchronization +# +# `zeCommandListHostSynchronize` and `zeCommandQueueSynchronize` block the calling thread +# until the work has completed, so no other task can run on it in the meantime. As CUDA.jl +# does, first busy-wait on a non-blocking query, which keeps the latency of short +# operations low, and then block in the driver on a separate thread, while the calling +# task waits for that thread without blocking the scheduler. + +export nonblocking_synchronize + +const SyncObject = Union{ZeImmediateCommandList, ZeCommandQueue} + +# with a zero timeout, a synchronization is a query +function check_done(res::ze_result_t) + if res == RESULT_NOT_READY + return false + elseif res == RESULT_SUCCESS + return true + else + throw_api_error(res) + end +end +Base.isdone(list::ZeImmediateCommandList) = + check_done(unchecked_zeCommandListHostSynchronize(list, 0)) +Base.isdone(queue::ZeCommandQueue) = + check_done(unchecked_zeCommandQueueSynchronize(queue, 0)) + +# the blocking synchronization, marked GC-safe so that it doesn't keep the GC from running +gcsafe_synchronize(list::ZeImmediateCommandList) = + @gcsafe_ccall libze_loader.zeCommandListHostSynchronize( + list::ze_command_list_handle_t, typemax(UInt64)::UInt64)::ze_result_t +gcsafe_synchronize(queue::ZeCommandQueue) = + @gcsafe_ccall libze_loader.zeCommandQueueSynchronize( + queue::ze_command_queue_handle_t, typemax(UInt64)::UInt64)::ze_result_t + + +## bidirectional channel + +# custom, unbuffered channel that supports returning a value to the sender +# without the need for a second channel +struct BidirectionalChannel{I,O} <: AbstractChannel{I} + cond_take::Threads.Condition # waiting for data to become available + cond_put::Threads.Condition # waiting for a writeable slot + cond_ret::Threads.Condition # waiting for a data to be returned + + function BidirectionalChannel{I,O}() where {I,O} + lock = ReentrantLock() + cond_put = Threads.Condition(lock) + cond_take = Threads.Condition(lock) + cond_ret = Threads.Condition(lock) + return new(cond_take, cond_put, cond_ret) + end +end + +Base.put!(c::BidirectionalChannel{I}, v) where {I} = put!(c, convert(I, v)) +function Base.put!(c::BidirectionalChannel{I,O}, v::I) where {I,O} + lock(c) + try + # wait for a slot to be available + while isempty(c.cond_take) + Base.wait(c.cond_put) + end + + # pass a value to the consumer + notify(c.cond_take, v, false, false) + + # wait for a return value to be produced + Base.wait(c.cond_ret)::O + finally + unlock(c) + end +end + +function Base.take!(f::Base.Callable, c::BidirectionalChannel{I,O}) where {I,O} + lock(c) + try + # notify the producer that we're ready to accept a value + notify(c.cond_put, nothing, false, false) + + # receive a value from the producer + v = Base.wait(c.cond_take)::I + + # return a value to the producer + ret = f(v)::O + notify(c.cond_ret, ret, false, false) + finally + unlock(c) + end +end + +Base.lock(c::BidirectionalChannel) = lock(c.cond_take) +Base.unlock(c::BidirectionalChannel) = unlock(c.cond_take) + + +## fast path + +# before blocking on a separate thread, which has some overhead, busy-wait on a query of +# the object to synchronize. when this returns true, the object still has to be +# synchronized, but that won't block anymore. +function spinning_synchronization(f, obj) + # fast path + f(obj) && return true + + # minimize latency of short operations by busy-waiting, + # initially without even yielding to other tasks + 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 + f(obj) && return true + spins += 1 + end + + return false +end + + +## slow path: synchronize on a separate thread + +const MAX_SYNC_THREADS = 4 +const sync_channels = Array{BidirectionalChannel{SyncObject,ze_result_t}}(undef, MAX_SYNC_THREADS) +const sync_channel_cursor = Threads.Atomic{UInt32}(1) +const sync_channel_lock = Base.ReentrantLock() + +function synchronization_worker(data) + i = Int(data) + chan = sync_channels[i] + + while true + # wait for work + take!(gcsafe_synchronize, chan) + end +end + +@noinline function create_synchronization_worker(i) + lock(sync_channel_lock) do + # test and test-and-set + if isassigned(sync_channels, i) + return + end + + # should be safe to assign before threads are running; + # any user will just submit work that makes it block + sync_channels[i] = BidirectionalChannel{SyncObject,ze_result_t}() + + # we don't know what the size of uv_thread_t is, so reserve enough space + tid = Ref{NTuple{32, UInt8}}(ntuple(i -> 0, 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) + end + + return +end + +""" + nonblocking_synchronize(list_or_queue) + +Wait for the work on an immediate command list or command queue to complete, like +[`synchronize`](@ref), but without blocking the Julia scheduler: other tasks keep running +while this one waits. +""" +function nonblocking_synchronize(obj::SyncObject) + if spinning_synchronization(Base.isdone, obj) + # done, so this doesn't block + synchronize(obj) + return + end + + # pick a worker channel: sticky per task, so repeated synchronizations from the + # same task always hit the same, already running worker thread. + tls = task_local_storage() + i = get!(tls, :ZeSyncChannel) do + mod1(Threads.atomic_add!(sync_channel_cursor, UInt32(1)), MAX_SYNC_THREADS) + end::Int + if !isassigned(sync_channels, i) + create_synchronization_worker(i) + end + chan = @inbounds sync_channels[i] + + # submit the object to synchronize; unlike with regular channels, this `put!` blocks + # until the worker has synchronized it and returned the result + res = put!(chan, obj) + res == RESULT_SUCCESS || throw_api_error(res) + + return +end diff --git a/lib/utils/APIUtils.jl b/lib/utils/APIUtils.jl index d8d30394..61ec0743 100644 --- a/lib/utils/APIUtils.jl +++ b/lib/utils/APIUtils.jl @@ -1,8 +1,8 @@ module APIUtils # helpers that facilitate working with C APIs -using GPUToolbox: @checked, @debug_ccall -export @checked, @debug_ccall +using GPUToolbox: @checked, @debug_ccall, @gcsafe_ccall +export @checked, @debug_ccall, @gcsafe_ccall include("enum.jl") end diff --git a/src/compiler/compilation.jl b/src/compiler/compilation.jl index 8f463103..7305e049 100644 --- a/src/compiler/compilation.jl +++ b/src/compiler/compilation.jl @@ -1,6 +1,9 @@ ## gpucompiler interface implementation -struct oneAPICompilerParams <: AbstractCompilerParams end +Base.@kwdef struct oneAPICompilerParams <: AbstractCompilerParams + sub_group_size::Union{Nothing,Int} = nothing +end + const oneAPICompilerConfig = CompilerConfig{SPIRVCompilerTarget, oneAPICompilerParams} const oneAPICompilerJob = CompilerJob{SPIRVCompilerTarget,oneAPICompilerParams} @@ -58,6 +61,11 @@ function GPUCompiler.finish_module!(job::oneAPICompilerJob, mod::LLVM.Module, Tuple{CompilerJob{SPIRVCompilerTarget}, typeof(mod), typeof(entry)}, job, mod, entry) + # Set the subgroup size + if job.config.params.sub_group_size !== nothing + metadata(entry)["intel_reqd_sub_group_size"] = MDNode([ConstantInt(Int32(job.config.params.sub_group_size))]) + end + # OpenCL 2.0 push!(metadata(mod)["opencl.ocl.version"], MDNode([ConstantInt(Int32(2)), @@ -273,11 +281,16 @@ function _driver_supports_bfloat16_spirv(dev=device()) end end -@noinline function _compiler_config(dev; kernel=true, name=nothing, always_inline=false, kwargs...) +@noinline function _compiler_config(dev; kernel=true, name=nothing, always_inline=false, sub_group_size=32, kwargs...) properties = oneL0.module_properties(dev) supports_fp16 = properties.fp16flags & oneL0.ZE_DEVICE_MODULE_FLAG_FP16 == oneL0.ZE_DEVICE_MODULE_FLAG_FP16 supports_fp64 = properties.fp64flags & oneL0.ZE_DEVICE_MODULE_FLAG_FP64 == oneL0.ZE_DEVICE_MODULE_FLAG_FP64 + if sub_group_size ∉ oneL0.compute_properties(dev).subGroupSizes + @error("$sub_group_size is not a valid sub-group size for this device.") + end + + # SPIR-V codegen path. The Aurora LTS NEO/IGC runtime only accepts SPIR-V from the # Khronos translator; the rolling stack uses the LLVM SPIR-V back-end. GPUCompiler picks # the tool from the target's `backend` field and loads the JLL lazily, so both can be @@ -311,7 +324,7 @@ end # create GPUCompiler objects target = SPIRVCompilerTarget(; backend, extensions = extensions_str, supports_fp16, supports_fp64, supports_bfloat16, driver = :intel, kwargs...) - params = oneAPICompilerParams() + params = oneAPICompilerParams(; sub_group_size) CompilerConfig(target, params; kernel, name, always_inline) end diff --git a/src/compiler/execution.jl b/src/compiler/execution.jl index 20742c2d..06749f94 100644 --- a/src/compiler/execution.jl +++ b/src/compiler/execution.jl @@ -4,7 +4,7 @@ export @oneapi, zefunction, kernel_convert ## high-level @oneapi interface const MACRO_KWARGS = [:launch] -const COMPILER_KWARGS = [:kernel, :name, :always_inline] +const COMPILER_KWARGS = [:kernel, :name, :always_inline, :sub_group_size] const LAUNCH_KWARGS = [:groups, :items, :queue] """ @@ -241,9 +241,9 @@ function launch_configuration(kernel::HostKernel{F,TT}) where {F,TT} # configurations, so roll our own version that behaves like CUDA's # occupancy API and assumes the kernel still does bounds checking. - kernel_props = oneL0.properties(kernel.fun) - group_size = if kernel_props.maxGroupSize !== missing - kernel_props.maxGroupSize + max_group_size = oneL0.max_group_size(kernel.fun) + group_size = if max_group_size !== missing + max_group_size else # without the MAX_GROUP_SIZE extension, we need to be conservative dev = kernel.fun.mod.device @@ -261,7 +261,7 @@ function launch_configuration(kernel::HostKernel{F,TT}) where {F,TT} # size but does not fold it into `maxGroupSize`, so account for it here. Rounded down to # a power of two, both because group sizes want to be anyway and to stay clear of the # limit rather than right at it. - spill = kernel_props.spillMemSize + spill = oneL0.spill_mem_size(kernel.fun) if spill > 0 && group_size * spill > MAX_GROUP_SCRATCH group_size = max(1, prevpow(2, max(1, MAX_GROUP_SCRATCH ÷ spill))) end @@ -360,7 +360,8 @@ end spill > s.scratch_hwm && scratch_hedge!(s, spill) append_launch!(s.list, kernel, groups) - oneL0.sync_each_submission() && oneL0.synchronize(s.list) + # wait cooperatively, as `synchronize` does, or the workaround blocks the thread + oneL0.sync_each_submission() && oneL0.nonblocking_synchronize(s.list) return end diff --git a/src/context.jl b/src/context.jl index 409a9d7a..6c9e188f 100644 --- a/src/context.jl +++ b/src/context.jl @@ -378,12 +378,15 @@ function synchronize_all_streams(ctx::ZeContext, dev::Union{ZeDevice, Nothing}) end """ - synchronize() - synchronize(stream::oneStream) + synchronize(; blocking=false) + synchronize(stream::oneStream; blocking=false) -Block the host thread until all operations on the calling task's stream for the current -context and device have completed: work appended to the immediate command list as well -as oneMKL work on the companion queue. +Block the calling task until all operations on its stream for the current context and +device have completed: work appended to the immediate command list as well as oneMKL work +on the companion queue. + +Unless `blocking` is set, other tasks keep running while waiting: the host thread is only +blocked in the driver when the work is already done. This is useful for timing operations or ensuring that GPU work has finished before accessing results on the CPU. @@ -398,18 +401,19 @@ println("GPU work completed") See also: [`global_stream`](@ref), [`context`](@ref), [`device`](@ref) """ -function oneL0.synchronize(s::oneStream) - oneL0.synchronize(s.list) +function oneL0.synchronize(s::oneStream; blocking::Bool=false) + sync = blocking ? oneL0.synchronize : oneL0.nonblocking_synchronize + sync(s.list) q = s.queue if q !== nothing - oneL0.synchronize(q) + sync(q) s.mkl_dirty = false end return end -function oneL0.synchronize() - oneL0.synchronize(global_stream(context(), device())) +function oneL0.synchronize(; blocking::Bool=false) + oneL0.synchronize(global_stream(context(), device()); blocking) end # Julia → MKL ordering: everything Julia appended to the task's immediate list must be diff --git a/src/oneAPI.jl b/src/oneAPI.jl index 6b6db42c..66b42a2a 100644 --- a/src/oneAPI.jl +++ b/src/oneAPI.jl @@ -13,6 +13,7 @@ using SpecialFunctions import Preferences import KernelAbstractions: KernelAbstractions +import KernelInterface using LLVM using LLVM.Interop @@ -77,12 +78,18 @@ include("gpuarrays.jl") include("random.jl") include("utils.jl") -include("oneAPIKernels.jl") +# KernelAbstractions +include("oneAPIKernelsOld.jl") import .oneAPIKernels: oneAPIBackend +export oneAPIBackend + +# KernelInterface +include("oneAPIKernels.jl") +import .oneAPIInterface + include("accumulate.jl") include("sorting.jl") include("indexing.jl") -export oneAPIBackend # precompilation workload (warms up the SPIR-V compilation pipeline) include("compiler/precompile.jl") diff --git a/src/oneAPIKernels.jl b/src/oneAPIKernels.jl index 7b90d2ca..7346df25 100644 --- a/src/oneAPIKernels.jl +++ b/src/oneAPIKernels.jl @@ -1,229 +1,191 @@ -module oneAPIKernels +module oneAPIInterface using ..oneAPI -using ..oneAPI: @device_override, SPIRVIntrinsics, method_table +using ..oneAPI: @device_override, SPIRVIntrinsics, method_table, kernel_convert, zefunction -import KernelAbstractions as KA +import KernelInterface as KI import StaticArrays -import Adapt - - ## Back-end Definition export oneAPIBackend -struct oneAPIBackend <: KA.GPU +struct oneAPIBackend <: KI.Backend prefer_blocks::Bool always_inline::Bool end -oneAPIBackend(; prefer_blocks = false, always_inline = false) = oneAPIBackend(prefer_blocks, always_inline) +KI.versioninfo(io::IO, ::oneAPIBackend) = oneAPI.versioninfo(io) -@inline KA.allocate(::oneAPIBackend, ::Type{T}, dims::Tuple; unified::Bool = false) where {T} = oneArray{T, length(dims), unified ? oneAPI.oneL0.SharedBuffer : oneAPI.oneL0.DeviceBuffer}(undef, dims) -@inline KA.zeros(::oneAPIBackend, ::Type{T}, dims::Tuple; unified::Bool = false) where {T} = fill!(oneArray{T, length(dims), unified ? oneAPI.oneL0.SharedBuffer : oneAPI.oneL0.DeviceBuffer}(undef, dims), zero(T)) -@inline KA.ones(::oneAPIBackend, ::Type{T}, dims::Tuple; unified::Bool = false) where {T} = fill!(oneArray{T, length(dims), unified ? oneAPI.oneL0.SharedBuffer : oneAPI.oneL0.DeviceBuffer}(undef, dims), one(T)) +oneAPIBackend(; prefer_blocks = false, always_inline = false) = oneAPIBackend(prefer_blocks, always_inline) -KA.get_backend(::oneArray) = oneAPIBackend() -# TODO should be non-blocking -KA.synchronize(::oneAPIBackend) = oneAPI.oneL0.synchronize() -KA.supports_float64(::oneAPIBackend) = false # TODO: Check if this is device dependent -KA.supports_unified(::oneAPIBackend) = true +@inline KI.allocate(::oneAPIBackend, ::Type{T}, dims::Tuple; unified::Bool = false) where {T} = oneArray{T, length(dims), unified ? oneAPI.oneL0.SharedBuffer : oneAPI.oneL0.DeviceBuffer}(undef, dims) -KA.functional(::oneAPIBackend) = oneAPI.functional() +KI.get_backend(::oneArray) = oneAPIBackend() +KI.synchronize(::oneAPIBackend) = oneAPI.oneL0.synchronize() +KI.supports_float64(::oneAPIBackend) = device_limits().supports_float64 +KI.supports_unified(::oneAPIBackend) = true +KI.supports_atomics(::oneAPIBackend) = true -Adapt.adapt_storage(::oneAPIBackend, a::AbstractArray) = Adapt.adapt(oneArray, a) -Adapt.adapt_storage(::oneAPIBackend, a::oneArray) = a -Adapt.adapt_storage(::KA.CPU, a::oneArray) = convert(Array, a) +KI.functional(::oneAPIBackend) = oneAPI.functional() # sparse arrays (oneMKL is only available on Linux) @static if Sys.islinux() - import GPUArrays, SparseArrays - KA.get_backend(::oneAPI.oneMKL.oneAbstractSparseMatrix) = oneAPIBackend() - # without this, `adapt_storage(::oneAPIBackend, ::AbstractArray)` would densify sparse arrays - Adapt.adapt_storage(::oneAPIBackend, a::GPUArrays.AbstractGPUSparseArray) = a - Adapt.adapt_storage(::KA.CPU, a::oneAPI.oneMKL.oneAbstractSparseMatrix) = SparseArrays.SparseMatrixCSC(a) + KI.get_backend(::oneAPI.oneMKL.oneAbstractSparseMatrix) = oneAPIBackend() end ## Memory Operations -function KA.copyto!(::oneAPIBackend, A, B) +function KI.copyto!(::oneAPIBackend, A, B) + length(A) == length(B) || + throw(ArgumentError("Arrays must have the same length, got $(length(A)) and $(length(B))")) copyto!(A, B) # TODO: Address device to host copies in jl being synchronizing + return A end +KI.unsafe_free!(A::oneArray) = oneAPI.unsafe_free!(A) + ## Device Operations -function KA.ndevices(::oneAPIBackend) +function KI.ndevices(::oneAPIBackend) return length(oneAPI.devices()) end -function KA.device(::oneAPIBackend)::Int - dev = oneAPI.device() +function device_index(dev)::Int devs = oneAPI.devices() idx = findfirst(==(dev), devs) return idx === nothing ? 1 : idx end +KI.device(::oneAPIBackend)::Int = device_index(oneAPI.device()) +KI.device(::oneAPIBackend, A::oneArray)::Int = device_index(oneAPI.device(A)) -function KA.device!(backend::oneAPIBackend, id::Int) - return oneAPI.device!(id) +function KI.device!(backend::oneAPIBackend, id::Int) + oneAPI.device!(id) + return end ## Kernel Launch -function KA.mkcontext(kernel::KA.Kernel{oneAPIBackend}, _ndrange, iterspace) - KA.CompilerMetadata{KA.ndrange(kernel), KA.DynamicCheck}(_ndrange, iterspace) -end -function KA.mkcontext(kernel::KA.Kernel{oneAPIBackend}, I, _ndrange, iterspace, - ::Dynamic) where Dynamic - KA.CompilerMetadata{KA.ndrange(kernel), Dynamic}(I, _ndrange, iterspace) -end - -function KA.launch_config(kernel::KA.Kernel{oneAPIBackend}, ndrange, workgroupsize) - if ndrange isa Integer - ndrange = (ndrange,) - end - if workgroupsize isa Integer - workgroupsize = (workgroupsize, ) - end +KI.argconvert(::oneAPIBackend, arg) = kernel_convert(arg) - # 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 - # (clamped to 1, since an empty ndrange cannot serve as a workgroup size) - KA.partition(kernel, ndrange, max.(ndrange, 1)) +function KI.kernel_function(backend::oneAPIBackend, f::F, tt::TT=Tuple{}; name = nothing, kwargs...) where {F,TT} + # compile for the sub-group width that `KI.sub_group_size` promises + sub_group_size = KI.sub_group_size(backend) + kern = if sub_group_size > 0 + zefunction(f, tt; name, backend.always_inline, sub_group_size, kwargs...) else - KA.partition(kernel, ndrange, workgroupsize) + zefunction(f, tt; name, backend.always_inline, kwargs...) end - - return ndrange, workgroupsize, iterspace, dynamic -end - -function threads_to_workgroupsize(threads, ndrange) - total = 1 - return map(ndrange) do n - x = max(1, min(div(threads, total), n)) - total *= x - return x + KI.Kernel{oneAPIBackend, typeof(kern)}(backend, kern) +end + +function KI.launch(obj::KI.Kernel{oneAPIBackend}, groups::Dims{3}, items::Dims{3}, args::Vararg{Any, N}; kwargs...) where {N} + # kernels are compiled for a device, and launched on the task's stream of the active one + obj.kern.fun.mod.device == device() || + throw(ArgumentError("Cannot launch a kernel compiled for another device than the active one")) + obj.kern(args...; items, groups, kwargs...) + return +end + +function KI.max_work_group_size(kernel::KI.Kernel{oneAPIBackend})::Int + fun = kernel.kern.fun + max_group_size = oneAPI.oneL0.max_group_size(fun) + # without the MAX_GROUP_SIZE extension, the device limit is all we know + return coalesce(max_group_size, device_limits(fun.mod.device).max_work_group_size) +end +function KI.launch_configuration( + kernel::KI.Kernel{oneAPIBackend}; nitems::Union{Integer, Nothing} = nothing, + max_work_group_size::Integer = typemax(Int) + ) + group_size = oneAPI.launch_configuration(kernel.kern) + return (; workgroupsize = Int(min(group_size, max_work_group_size))) +end +# querying the device allocates, so cache what every launch needs +const DeviceLimits = @NamedTuple{ + max_work_group_size::Int, max_work_group_dims::NTuple{3, Int}, max_num_groups::NTuple{3, Int}, + sub_group_size::Int, supports_float16::Bool, supports_float64::Bool, +} +function device_limits(dev::oneAPI.oneL0.ZeDevice = device()) + limits = get!(task_local_storage(), :oneAPIDeviceLimits) do + Dict{oneAPI.oneL0.ZeDevice, DeviceLimits}() + end::Dict{oneAPI.oneL0.ZeDevice, DeviceLimits} + get!(limits, dev) do + props = oneAPI.oneL0.compute_properties(dev) + module_props = oneAPI.oneL0.module_properties(dev) + # the sub-group width that `kernel_function` compiles for: the width `@oneapi` defaults + # to if the device supports it, and 0 if the device has no sub-groups + sg_sizes = props.subGroupSizes + sub_group_size = 32 in sg_sizes ? 32 : maximum(sg_sizes; init = 0) + (; max_work_group_size = props.maxTotalGroupSize, + max_work_group_dims = (props.maxGroupSizeX, props.maxGroupSizeY, props.maxGroupSizeZ), + max_num_groups = (props.maxGroupCountX, props.maxGroupCountY, props.maxGroupCountZ), + sub_group_size, + supports_float16 = module_props.flags & oneAPI.oneL0.ZE_DEVICE_MODULE_FLAG_FP16 != 0, + supports_float64 = module_props.flags & oneAPI.oneL0.ZE_DEVICE_MODULE_FLAG_FP64 != 0) end end - -function (obj::KA.Kernel{oneAPIBackend})(args...; ndrange=nothing, workgroupsize=nothing) - backend = KA.backend(obj) - - 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) - - # If the kernel is statically sized we can tell the compiler about that - if KA.workgroupsize(obj) <: KA.StaticSize - # TODO: maxthreads - # maxthreads = prod(KA.get(KA.workgroupsize(obj))) - else - # maxthreads = nothing - end - - kernel = @oneapi launch = false always_inline = backend.always_inline obj.f(ctx, args...) - - # figure out the optimal workgroupsize automatically - if KA.workgroupsize(obj) <: KA.DynamicSize && workgroupsize === nothing - items = oneAPI.launch_configuration(kernel) - - if backend.prefer_blocks - # Prefer blocks over threads: - # Reducing the workgroup size (items) increases the number of workgroups (blocks). - # We use a simple heuristic here since we lack full occupancy info (max_blocks) from launch_configuration. - - # If the total range is large enough, full workgroups are fine. - # If the range is small, we might want to reduce 'items' to create more blocks to fill the GPU. - # (Simplified logic compared to CUDA.jl which uses explicit occupancy calculators) - total_items = prod(ndrange) - if total_items < items * 16 # Heuristic factor - # Force at least a few blocks if possible by reducing items per block - target_blocks = 16 # Target at least 16 blocks - items = max(1, min(items, cld(total_items, target_blocks))) - end - end - - workgroupsize = threads_to_workgroupsize(items, ndrange) - iterspace, dynamic = KA.partition(obj, ndrange, workgroupsize) - 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 - kernel(ctx, args...; items, groups) - - return nothing +KI.max_work_group_size(::oneAPIBackend)::Int = device_limits().max_work_group_size +KI.max_work_group_dims(::oneAPIBackend)::NTuple{3, Int} = device_limits().max_work_group_dims +KI.max_num_groups(::oneAPIBackend)::NTuple{3, Int} = device_limits().max_num_groups +KI.sub_group_size(::oneAPIBackend)::Int = device_limits().sub_group_size +function KI.multiprocessor_count(::oneAPIBackend)::Int + oneAPI.oneL0.properties(device()).numSlices end +KI.supports_subgroups(::oneAPIBackend) = device_limits().sub_group_size > 0 +function KI.supports_shuffle(::oneAPIBackend, ::Type{T}) where {T} + T in SPIRVIntrinsics.gentypes || return false + T === Float64 && return device_limits().supports_float64 + T === Float16 && return device_limits().supports_float16 + return true +end ## Indexing Functions +## COV_EXCL_START -@device_override @inline function KA.__index_Local_Linear(ctx) - return get_local_id() -end +# computed with `% T`, which unlike `T(x)` has no error path -@device_override @inline function KA.__index_Group_Linear(ctx) - return get_group_id() +@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() +@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()] +@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()] +@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(), get_local_id()) -end +@device_override KI.get_sub_group_size(::Type{T}) where {T} = get_sub_group_size() % T -@device_override @inline function KA.__validindex(ctx) - if KA.__dynamic_checkbounds(ctx) - I = @inbounds KA.expand(KA.__iterspace(ctx), get_group_id(), get_local_id()) - return I in KA.__ndrange(ctx) - else - return true - end -end +@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 = oneAPI.emit_localmemory(T, Val(prod(Dims))) oneDeviceArray(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() # Fence both local and global memory across the workgroup barrier, matching CUDA # `__syncthreads` semantics. `barrier(0)` lowers to `OpControlBarrier` with # `SequentiallyConsistent` but WITHOUT any storage-class bit, which the SPIR-V spec @@ -233,18 +195,23 @@ end barrier(SPIRVIntrinsics.LOCAL_MEM_FENCE | SPIRVIntrinsics.GLOBAL_MEM_FENCE) end -@device_override @inline function KA.__print(args...) - oneAPI._print(args...) +@device_override @inline function KI.sub_group_barrier() + sub_group_barrier(SPIRVIntrinsics.LOCAL_MEM_FENCE | SPIRVIntrinsics.GLOBAL_MEM_FENCE) end +@device_override function KI.shfl_down(val::T, offset::Integer) where T + sub_group_shuffle(val, get_sub_group_local_id() + offset) +end -## Other +@device_override @inline function KI._print(args...) + oneAPI._print(args...) +end -Adapt.adapt_storage(to::KA.ConstAdaptor, a::oneDeviceArray) = Base.Experimental.Const(a) +## COV_EXCL_STOP -KA.argconvert(::KA.Kernel{oneAPIBackend}, arg) = kernel_convert(arg) +## Other -function KA.priority!(::oneAPIBackend, prio::Symbol) +function KI.priority!(::oneAPIBackend, prio::Symbol) if !(prio in (:high, :normal, :low)) error("priority must be one of :high, :normal, :low") end diff --git a/src/oneAPIKernelsOld.jl b/src/oneAPIKernelsOld.jl new file mode 100644 index 00000000..854e8ab4 --- /dev/null +++ b/src/oneAPIKernelsOld.jl @@ -0,0 +1,282 @@ +module oneAPIKernels + +using ..oneAPI +using ..oneAPI: @device_override, SPIRVIntrinsics, method_table + +import KernelAbstractions as KA + +import StaticArrays + +import Adapt + + +## Back-end Definition + +export oneAPIBackend + +struct oneAPIBackend <: KA.GPU + prefer_blocks::Bool + always_inline::Bool +end + +oneAPIBackend(; prefer_blocks = false, always_inline = false) = oneAPIBackend(prefer_blocks, always_inline) + +@inline KA.allocate(::oneAPIBackend, ::Type{T}, dims::Tuple; unified::Bool = false) where {T} = oneArray{T, length(dims), unified ? oneAPI.oneL0.SharedBuffer : oneAPI.oneL0.DeviceBuffer}(undef, dims) +@inline KA.zeros(::oneAPIBackend, ::Type{T}, dims::Tuple; unified::Bool = false) where {T} = fill!(oneArray{T, length(dims), unified ? oneAPI.oneL0.SharedBuffer : oneAPI.oneL0.DeviceBuffer}(undef, dims), zero(T)) +@inline KA.ones(::oneAPIBackend, ::Type{T}, dims::Tuple; unified::Bool = false) where {T} = fill!(oneArray{T, length(dims), unified ? oneAPI.oneL0.SharedBuffer : oneAPI.oneL0.DeviceBuffer}(undef, dims), one(T)) + +KA.get_backend(::oneArray) = oneAPIBackend() +# TODO should be non-blocking +KA.synchronize(::oneAPIBackend) = oneAPI.oneL0.synchronize() +KA.supports_float64(::oneAPIBackend) = false # TODO: Check if this is device dependent +KA.supports_unified(::oneAPIBackend) = true + +KA.functional(::oneAPIBackend) = oneAPI.functional() + +Adapt.adapt_storage(::oneAPIBackend, a::AbstractArray) = Adapt.adapt(oneArray, a) +Adapt.adapt_storage(::oneAPIBackend, a::oneArray) = a +Adapt.adapt_storage(::KA.CPU, a::oneArray) = convert(Array, a) + +# sparse arrays (oneMKL is only available on Linux) +@static if Sys.islinux() + import GPUArrays, SparseArrays + KA.get_backend(::oneAPI.oneMKL.oneAbstractSparseMatrix) = oneAPIBackend() + # without this, `adapt_storage(::oneAPIBackend, ::AbstractArray)` would densify sparse arrays + Adapt.adapt_storage(::oneAPIBackend, a::GPUArrays.AbstractGPUSparseArray) = a + Adapt.adapt_storage(::KA.CPU, a::oneAPI.oneMKL.oneAbstractSparseMatrix) = SparseArrays.SparseMatrixCSC(a) +end + +## Memory Operations + +function KA.copyto!(::oneAPIBackend, A, B) + copyto!(A, B) + # TODO: Address device to host copies in jl being synchronizing +end + + +## Device Operations + +function KA.ndevices(::oneAPIBackend) + return length(oneAPI.devices()) +end + +function KA.device(::oneAPIBackend)::Int + dev = oneAPI.device() + devs = oneAPI.devices() + idx = findfirst(==(dev), devs) + return idx === nothing ? 1 : idx +end + +function KA.device!(backend::oneAPIBackend, id::Int) + return oneAPI.device!(id) +end + + +## Kernel Launch + +function KA.mkcontext(kernel::KA.Kernel{oneAPIBackend}, _ndrange, iterspace) + KA.CompilerMetadata{KA.ndrange(kernel), KA.DynamicCheck}(_ndrange, iterspace) +end +function KA.mkcontext(kernel::KA.Kernel{oneAPIBackend}, I, _ndrange, iterspace, + ::Dynamic) where Dynamic + KA.CompilerMetadata{KA.ndrange(kernel), Dynamic}(I, _ndrange, iterspace) +end + +function KA.launch_config(kernel::KA.Kernel{oneAPIBackend}, 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 + # (clamped to 1, since an empty ndrange cannot serve as a workgroup size) + KA.partition(kernel, ndrange, max.(ndrange, 1)) + 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 = max(1, min(div(threads, total), n)) + total *= x + return x + end +end + +function (obj::KA.Kernel{oneAPIBackend})(args...; ndrange=nothing, workgroupsize=nothing) + backend = KA.backend(obj) + + 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) + + # If the kernel is statically sized we can tell the compiler about that + if KA.workgroupsize(obj) <: KA.StaticSize + # TODO: maxthreads + # maxthreads = prod(KA.get(KA.workgroupsize(obj))) + else + # maxthreads = nothing + end + + kernel = @oneapi launch = false always_inline = backend.always_inline obj.f(ctx, args...) + + # figure out the optimal workgroupsize automatically + if KA.workgroupsize(obj) <: KA.DynamicSize && workgroupsize === nothing + items = oneAPI.launch_configuration(kernel) + + if backend.prefer_blocks + # Prefer blocks over threads: + # Reducing the workgroup size (items) increases the number of workgroups (blocks). + # We use a simple heuristic here since we lack full occupancy info (max_blocks) from launch_configuration. + + # If the total range is large enough, full workgroups are fine. + # If the range is small, we might want to reduce 'items' to create more blocks to fill the GPU. + # (Simplified logic compared to CUDA.jl which uses explicit occupancy calculators) + total_items = prod(ndrange) + if total_items < items * 16 # Heuristic factor + # Force at least a few blocks if possible by reducing items per block + target_blocks = 16 # Target at least 16 blocks + items = max(1, min(items, cld(total_items, target_blocks))) + end + end + + workgroupsize = threads_to_workgroupsize(items, ndrange) + iterspace, dynamic = KA.partition(obj, ndrange, workgroupsize) + 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 + kernel(ctx, args...; items, groups) + + return nothing +end + + +## Indexing Functions + +@device_override @inline function KA.__index_Local_Linear(ctx) + return get_local_id() +end + +@device_override @inline function KA.__index_Group_Linear(ctx) + return get_group_id() +end + +@device_override @inline function KA.__index_Global_Linear(ctx) + return get_global_id() +end + +@device_override @inline function KA.__index_Local_Cartesian(ctx) + @inbounds KA.workitems(KA.__iterspace(ctx))[get_local_id()] +end + +@device_override @inline function KA.__index_Group_Cartesian(ctx) + @inbounds KA.blocks(KA.__iterspace(ctx))[get_group_id()] +end + +@device_override @inline function KA.__index_Global_Cartesian(ctx) + return @inbounds KA.expand(KA.__iterspace(ctx), get_group_id(), get_local_id()) +end + +@device_override @inline function KA.__validindex(ctx) + if KA.__dynamic_checkbounds(ctx) + I = @inbounds KA.expand(KA.__iterspace(ctx), get_group_id(), get_local_id()) + 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 = oneAPI.emit_localmemory(T, Val(prod(Dims))) + oneDeviceArray(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() + # Fence both local and global memory across the workgroup barrier, matching CUDA + # `__syncthreads` semantics. `barrier(0)` lowers to `OpControlBarrier` with + # `SequentiallyConsistent` but WITHOUT any storage-class bit, which the SPIR-V spec + # treats as ordering *no* memory — so shared-local or global writes are not guaranteed + # visible to other work-items after the barrier. `LOCAL_MEM_FENCE | GLOBAL_MEM_FENCE` + # ORs in the WorkgroupMemory/CrossWorkgroupMemory fence bits. + barrier(SPIRVIntrinsics.LOCAL_MEM_FENCE | SPIRVIntrinsics.GLOBAL_MEM_FENCE) +end + +@device_override @inline function KA.__print(args...) + oneAPI._print(args...) +end + + +## Other + +Adapt.adapt_storage(to::KA.ConstAdaptor, a::oneDeviceArray) = Base.Experimental.Const(a) + +KA.argconvert(::KA.Kernel{oneAPIBackend}, arg) = kernel_convert(arg) + +function KA.priority!(::oneAPIBackend, prio::Symbol) + if !(prio in (:high, :normal, :low)) + error("priority must be one of :high, :normal, :low") + end + + priority_enum = if prio == :high + oneAPI.oneL0.ZE_COMMAND_QUEUE_PRIORITY_PRIORITY_HIGH + elseif prio == :low + oneAPI.oneL0.ZE_COMMAND_QUEUE_PRIORITY_PRIORITY_LOW + else + oneAPI.oneL0.ZE_COMMAND_QUEUE_PRIORITY_NORMAL + end + + ctx = oneAPI.context() + dev = oneAPI.device() + + # drain the task's current stream before swapping it out, so operations submitted + # to the new stream cannot overtake in-flight work on the old one + oneAPI.oneL0.synchronize(oneAPI.global_stream(ctx, dev)) + + # Replace the stream in task_local_storage. `create_stream` registers the + # replacement so `synchronize_all_streams`/`release` can drain it before freeing a + # buffer whose in-flight work it references; otherwise all work after `priority!` + # runs on an unregistered stream and a freed buffer can be reused while its kernel + # is still running (use-after-free → banned context on the LTS NEO stack). The old + # stream stays registered until its task dies, like replaced queues before it. + new_stream = oneAPI.create_stream(ctx, dev, priority_enum) + task_local_storage((:oneStream, ctx, dev), new_stream) + + # the cached SYCL queue wraps the old stream's companion queue; drop it so the next + # oneMKL call recreates it against the new stream (the old one was just drained) + delete!(task_local_storage(), (:SYCLQueue, ctx, dev)) + + return nothing +end + +end diff --git a/src/utils.jl b/src/utils.jl index ddf4b1f9..05bceb8a 100644 --- a/src/utils.jl +++ b/src/utils.jl @@ -27,11 +27,19 @@ function versioninfo(io::IO=stdout) println(io, "- LLVM: $(LLVM.version())") println(io) + + get_module(name::Symbol) = (name, getfield(oneAPI, name)) + function get_module(pkg::Tuple{String, String}) + id = Base.PkgId(Base.UUID(pkg[1]), pkg[2]) + (pkg[2], get(Base.loaded_modules, id, nothing)) + end + println(io, "Julia packages:") println(io, "- oneAPI.jl: $(Base.pkgversion(oneAPI))") - for name in [:GPUArrays, :GPUCompiler, :KernelAbstractions, :LLVM, :SPIRVIntrinsics] - mod = getfield(oneAPI, name) - println(io, "- $(name): $(Base.pkgversion(mod))") + for pkg in [:GPUArrays, :GPUCompiler, ("63c18a36-062a-441e-b654-da1e3ab1ce7c", "KernelAbstractions"), + :KernelInterface, :LLVM, :SPIRVIntrinsics] + name, mod = get_module(pkg) + isnothing(mod) || println(io, "- $(name): $(Base.pkgversion(mod))") end println(io) diff --git a/test/Project.toml b/test/Project.toml index 190beeab..e68d74ca 100644 --- a/test/Project.toml +++ b/test/Project.toml @@ -9,6 +9,7 @@ GPUArrays = "0c68f7d7-f131-5f86-a1c3-88cf8149b2d7" 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" NEO_jll = "700fe977-ac61-5f37-bbc8-c6c4b2b6a9fd" ParallelTestRunner = "d3525ed8-44d0-4b2c-a655-542cee43accc" @@ -25,5 +26,8 @@ libigc_jll = "94295238-5935-5bd7-bb0f-b00942e9bdd5" oneAPI = "8f75cd03-7ff8-4ecb-9b8f-daf728133b1b" oneAPI_Support_jll = "b049733a-a71d-5ed3-8eba-7d323ac00b36" +[sources] +KernelInterface = {url = "https://github.com/JuliaGPU/KernelAbstractions.jl", rev = "tb/ki-0.3", subdir = "lib/KernelInterface"} + [compat] ParallelTestRunner = "2.8" diff --git a/test/execution.jl b/test/execution.jl index 3efcdbf1..330a1612 100644 --- a/test/execution.jl +++ b/test/execution.jl @@ -743,6 +743,51 @@ end @test all(results) end +# burns `iters` dependent steps per work-item, which the compiler cannot fold away +function slow_kernel(a, iters) + i = get_global_id() + acc = i % UInt32 + for k in UInt32(1):iters + acc = acc * 0x0019660d + k + end + @inbounds a[i] = acc + return +end + +@testset "cooperative synchronize" begin + a = oneArray{UInt32}(undef, 64) + slow(iters) = @oneapi items=64 slow_kernel(a, UInt32(iters)) + slow(1) + synchronize() + # warm up the slow path of `synchronize`: compiling it would end the calibration early + slow(2^16) + synchronize() + + # make the kernel run for a while + iters = 2^16 + while iters < 2^30 && @elapsed((slow(iters); synchronize())) < 0.1 + iters *= 4 + end + + # other tasks keep running while one waits + stamps = UInt64[] + waiting = Ref(true) + ticker = @async while waiting[] + push!(stamps, time_ns()) + sleep(0.001) + end + slow(iters) + synchronize() + done = time_ns() + waiting[] = false + wait(ticker) + @test count(<(done), stamps) > 10 + + # blocking synchronization is still available + slow(1) + @test synchronize(; blocking=true) === nothing +end + ############################################################################################ # Keep allocation consumers at top level so kernels do not capture test state. diff --git a/test/kernelinterface.jl b/test/kernelinterface.jl new file mode 100644 index 00000000..75b704c4 --- /dev/null +++ b/test/kernelinterface.jl @@ -0,0 +1,6 @@ +import KernelInterface +using oneAPI.oneAPIInterface + +include(joinpath(dirname(pathof(KernelInterface)), "..", "test", "testsuite.jl")) + +Testsuite.testsuite(oneAPIInterface.oneAPIBackend(), oneArray)