From 7a1140be958a23be970699691e74f250f6c41704 Mon Sep 17 00:00:00 2001 From: Christian Guinard <28689358+christiangnrd@users.noreply.github.com> Date: Sat, 22 Aug 2026 15:32:34 -0300 Subject: [PATCH 1/9] Only list KernelAbstractions in versioninfo when it is loaded In preparation for KernelAbstractions becoming a weak dependency. --- src/utils.jl | 14 +++++++++++--- 1 file changed, 11 insertions(+), 3 deletions(-) diff --git a/src/utils.jl b/src/utils.jl index ddf4b1f9..03d4de12 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"), + :LLVM, :SPIRVIntrinsics] + name, mod = get_module(pkg) + isnothing(mod) || println(io, "- $(name): $(Base.pkgversion(mod))") end println(io) From 1aca77181310f360e2a3806e5edb5eac0775460d Mon Sep 17 00:00:00 2001 From: Christian Guinard <28689358+christiangnrd@users.noreply.github.com> Date: Tue, 6 Jan 2026 15:40:48 -0400 Subject: [PATCH 2/9] Allow setting the sub-group size of a kernel The `sub_group_size` compiler keyword sets the `intel_reqd_sub_group_size` metadata of the kernel, defaulting to 32. --- src/compiler/compilation.jl | 19 ++++++++++++++++--- src/compiler/execution.jl | 2 +- 2 files changed, 17 insertions(+), 4 deletions(-) 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..1f499ddf 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] """ From f9acd84121c2b4abd9d6c1a2b17a414438540c15 Mon Sep 17 00:00:00 2001 From: Christian Guinard <28689358+christiangnrd@users.noreply.github.com> Date: Sat, 22 Aug 2026 15:34:15 -0300 Subject: [PATCH 3/9] Add a KernelInterface back end Implement KernelInterface next to the KernelAbstractions back end, which moves to oneAPIKernelsOld.jl until the port to KernelAbstractions 0.10. The per-dimension launch limits (`max_work_group_dims`, `max_num_groups`) come from the Level Zero compute properties. Every auto-sized launch queries them, so they are cached per device. The KernelInterface test suite runs as the `kernelinterface` test. Its events test is skipped when every submission synchronizes (the Aurora LTS workaround, ONEAPI_SYNC_EACH_SUBMISSION), as the ordering it checks is not observable then. Co-authored-by: Tim Besard --- Project.toml | 2 + src/oneAPI.jl | 11 +- src/oneAPIKernels.jl | 248 +++++++++++++++-------------------- src/oneAPIKernelsOld.jl | 282 ++++++++++++++++++++++++++++++++++++++++ src/utils.jl | 2 +- test/Project.toml | 1 + test/kernelinterface.jl | 11 ++ 7 files changed, 409 insertions(+), 148 deletions(-) create mode 100644 src/oneAPIKernelsOld.jl create mode 100644 test/kernelinterface.jl diff --git a/Project.toml b/Project.toml index 88b57ba4..e21f1f71 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" @@ -42,6 +43,7 @@ GPUArrays = "11.5.14" GPUCompiler = "2.9" GPUToolbox = "3.1" KernelAbstractions = "0.9.39" +KernelInterface = "0.2.3" LLVM = "6, 7, 8, 9" NEO_jll = "=26.18.38308" PrecompileTools = "1" 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..f38be558 100644 --- a/src/oneAPIKernels.jl +++ b/src/oneAPIKernels.jl @@ -1,55 +1,46 @@ -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.GPU prefer_blocks::Bool always_inline::Bool end +KI.versioninfo(io::IO, ::oneAPIBackend) = oneAPI.versioninfo(io) + 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)) +@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) +@inline KI.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 KI.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() +KI.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 +KI.synchronize(::oneAPIBackend) = oneAPI.oneL0.synchronize() +KI.supports_float64(::oneAPIBackend) = false # TODO: Check if this is device dependent +KI.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) +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) copyto!(A, B) # TODO: Address device to host copies in jl being synchronizing end @@ -57,173 +48,135 @@ end ## Device Operations -function KA.ndevices(::oneAPIBackend) +function KI.ndevices(::oneAPIBackend) return length(oneAPI.devices()) end -function KA.device(::oneAPIBackend)::Int +function KI.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) +function KI.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 +KI.argconvert(::oneAPIBackend, arg) = kernel_convert(arg) -function KA.launch_config(kernel::KA.Kernel{oneAPIBackend}, ndrange, workgroupsize) - if ndrange isa Integer - ndrange = (ndrange,) - end - if workgroupsize isa Integer - workgroupsize = (workgroupsize, ) - end +function KI.kernel_function(::oneAPIBackend, f::F, tt::TT=Tuple{}; name = nothing, kwargs...) where {F,TT} + kern = zefunction(f, tt; name, kwargs...) + KI.Kernel{oneAPIBackend, typeof(kern)}(oneAPIBackend(), kern) +end - # partition checked that the ndrange's agreed - if KA.ndrange(kernel) <: KA.StaticSize - ndrange = nothing - end +function (obj::KI.Kernel{oneAPIBackend})(args...; numworkgroups=(), workgroupsize=(), ndrange=(), max_work_group_size=typemax(Int)) + KI.check_launch_args(numworkgroups, workgroupsize, ndrange) + prod(ndrange) == 0 && return nothing - 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 + numworkgroups, workgroupsize = KI.auto_launch_sizes(obj, numworkgroups, workgroupsize, ndrange, max_work_group_size) + items = (workgroupsize..., ntuple(_ -> 1, 3 - length(workgroupsize))...) + groups = (numworkgroups..., ntuple(_ -> 1, 3 - length(numworkgroups))...) - return ndrange, workgroupsize, iterspace, dynamic + obj.kern(args...; items, groups) + return nothing 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 +function KI.kernel_max_work_group_size(kernel::KI.Kernel{<:oneAPIBackend}; max_work_items::Int=typemax(Int))::Int + group_size = oneAPI.launch_configuration(kernel.kern) + Int(min(group_size, max_work_items)) +end +# querying the device allocates, so cache the limits that every auto-sized launch needs +const DeviceLimits = @NamedTuple{max_work_group_size::Int, max_work_group_dims::NTuple{3, Int}, + max_num_groups::NTuple{3, Int}} +function device_limits() + dev = device()::oneAPI.oneL0.ZeDevice + 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) + (; 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)) 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))) +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 +function KI.sub_group_size(::oneAPIBackend)::Int + sg_sizes = oneAPI.oneL0.compute_properties(device()).subGroupSizes + if 32 in sg_sizes + return 32 + elseif 64 in sg_sizes + return 64 + elseif 16 in sg_sizes + return 16 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) + return 1 end +end +function KI.multiprocessor_count(::oneAPIBackend)::Int + oneAPI.oneL0.properties(device()).numSlices +end - groups = length(KA.blocks(iterspace)) - items = length(KA.workitems(iterspace)) - - if groups == 0 - return nothing - end +function KI.shfl_down_types(::oneAPIBackend) + res = copy(SPIRVIntrinsics.gentypes) - # Launch kernel - kernel(ctx, args...; items, groups) + res = setdiff(res, [Float64]) - return nothing + return res end - ## Indexing Functions - -@device_override @inline function KA.__index_Local_Linear(ctx) - return get_local_id() +## COV_EXCL_START +@device_override @inline function KI.get_local_id(::Type{T}) where {T} + return (; x = T(get_local_id(1)), y = T(get_local_id(2)), z = T(get_local_id(3))) end -@device_override @inline function KA.__index_Group_Linear(ctx) - return get_group_id() +@device_override @inline function KI.get_group_id(::Type{T}) where {T} + return (; x = T(get_group_id(1)), y = T(get_group_id(2)), z = T(get_group_id(3))) end -@device_override @inline function KA.__index_Global_Linear(ctx) - return get_global_id() +@device_override @inline function KI.get_global_id(::Type{T}) where {T} + return (; x = T(get_global_id(1)), y = T(get_global_id(2)), z = T(get_global_id(3))) 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 = T(get_local_size(1)), y = T(get_local_size(2)), z = T(get_local_size(3))) 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 = T(get_num_groups(1)), y = T(get_num_groups(2)), z = T(get_num_groups(3))) end -@device_override @inline function KA.__index_Global_Cartesian(ctx) - return @inbounds KA.expand(KA.__iterspace(ctx), get_group_id(), get_local_id()) +@device_override @inline function KI.get_global_size(::Type{T}) where {T} + return (; x = T(get_global_size(1)), y = T(get_global_size(2)), z = T(get_global_size(3))) 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 +@device_override KI.get_sub_group_size() = get_sub_group_size() % UInt32 + +@device_override KI.get_max_sub_group_size() = get_max_sub_group_size() % UInt32 + +@device_override KI.get_num_sub_groups() = get_num_sub_groups() % UInt32 +@device_override KI.get_sub_group_id() = get_sub_group_id() % UInt32 + +@device_override KI.get_sub_group_local_id() = get_sub_group_local_id() % UInt32 ## 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 +186,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 03d4de12..05bceb8a 100644 --- a/src/utils.jl +++ b/src/utils.jl @@ -37,7 +37,7 @@ function versioninfo(io::IO=stdout) println(io, "Julia packages:") println(io, "- oneAPI.jl: $(Base.pkgversion(oneAPI))") for pkg in [:GPUArrays, :GPUCompiler, ("63c18a36-062a-441e-b654-da1e3ab1ce7c", "KernelAbstractions"), - :LLVM, :SPIRVIntrinsics] + :KernelInterface, :LLVM, :SPIRVIntrinsics] name, mod = get_module(pkg) isnothing(mod) || println(io, "- $(name): $(Base.pkgversion(mod))") end diff --git a/test/Project.toml b/test/Project.toml index 190beeab..fe34ca49 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" diff --git a/test/kernelinterface.jl b/test/kernelinterface.jl new file mode 100644 index 00000000..0854bcfe --- /dev/null +++ b/test/kernelinterface.jl @@ -0,0 +1,11 @@ +import KernelInterface +using oneAPI.oneAPIInterface + +include(joinpath(dirname(pathof(KernelInterface)), "..", "test", "testsuite.jl")) + +skip_tests = Set{String}() +# the events test checks that a waiting task blocks on work queued elsewhere, which is not +# observable when every submission synchronizes (the Aurora LTS workaround) +oneAPI.oneL0.sync_each_submission() && push!(skip_tests, "Events") + +Testsuite.testsuite(oneAPIInterface.oneAPIBackend, "oneAPI", oneAPI, oneArray, oneAPI.oneDeviceArray; skip_tests) From 116ace95495aaabc5f9b368b02edb67f2bb1b35a Mon Sep 17 00:00:00 2001 From: Tim Besard Date: Mon, 28 Sep 2026 10:12:28 +0200 Subject: [PATCH 4/9] Cache a kernel's maximum group size Like the spill size, cache the kernel's `maxGroupSize` in the `ZeKernel`, so that `launch_configuration` (and the KernelInterface launch validation, which checks every explicit work-group size against it) don't query the kernel properties, which allocates, on every launch. --- lib/level-zero/module.jl | 15 ++++++++++++++- src/compiler/execution.jl | 8 ++++---- 2 files changed, 18 insertions(+), 5 deletions(-) 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/src/compiler/execution.jl b/src/compiler/execution.jl index 1f499ddf..ed390744 100644 --- a/src/compiler/execution.jl +++ b/src/compiler/execution.jl @@ -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 From c8c76613959b9f709c3dedb5f98ec75b487c26df Mon Sep 17 00:00:00 2001 From: Tim Besard Date: Mon, 28 Sep 2026 10:12:28 +0200 Subject: [PATCH 5/9] Port the KernelInterface back end to KernelInterface 0.3 - subtype `KI.Backend`, and implement `KI.launch` instead of the call method: KernelInterface now validates the launch geometry itself; - implement the four primitive index queries with `% T`, and let KernelInterface derive the global ones; the sub-group queries take a type; - `max_work_group_size(kernel)` is the kernel's legal maximum, and `launch_configuration` the recommended size; - `supports_subgroups`/`supports_shuffle` replace `shfl_down_types`, and `supports_float64` and shuffles of `Float16`/`Float64` depend on the active device. `kernel_function` compiles for the sub-group width that `sub_group_size` reports, as KernelInterface now guarantees; - `kernel_function` keeps the backend it was given, so `always_inline` applies; - `copyto!` checks the lengths and returns the destination, `device!` returns `nothing`, and `device(backend, A)` and `unsafe_free!` are implemented; - launching a kernel on another device than it was compiled for throws; - drop `zeros`/`ones`, which KernelInterface implements generically. The testsuite now skips its events test by itself for back ends that record events by synchronizing, as this one does, so the explicit skip is gone. The typed index queries need SPIRVIntrinsics 1.1.3, which truncates the 3-D built-ins without producing illegal vector types. KernelInterface 0.3 isn't registered yet, so take it from its branch with `[sources]`, and develop it on Julia 1.10, which ignores `[sources]`. --- .buildkite/pipeline.yml | 5 ++ .github/workflows/docs.yml | 6 ++ .gitlab-ci.yml | 10 +++ Project.toml | 7 +- src/oneAPIKernels.jl | 135 +++++++++++++++++++------------------ test/Project.toml | 3 + test/kernelinterface.jl | 7 +- 7 files changed, 101 insertions(+), 72 deletions(-) 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 e21f1f71..3c008c85 100644 --- a/Project.toml +++ b/Project.toml @@ -33,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" @@ -43,12 +46,12 @@ GPUArrays = "11.5.14" GPUCompiler = "2.9" GPUToolbox = "3.1" KernelAbstractions = "0.9.39" -KernelInterface = "0.2.3" +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/src/oneAPIKernels.jl b/src/oneAPIKernels.jl index f38be558..1ed41b55 100644 --- a/src/oneAPIKernels.jl +++ b/src/oneAPIKernels.jl @@ -11,7 +11,7 @@ import StaticArrays export oneAPIBackend -struct oneAPIBackend <: KI.GPU +struct oneAPIBackend <: KI.Backend prefer_blocks::Bool always_inline::Bool end @@ -21,14 +21,13 @@ KI.versioninfo(io::IO, ::oneAPIBackend) = oneAPI.versioninfo(io) oneAPIBackend(; prefer_blocks = false, always_inline = false) = oneAPIBackend(prefer_blocks, always_inline) @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) -@inline KI.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 KI.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)) KI.get_backend(::oneArray) = oneAPIBackend() # TODO should be non-blocking KI.synchronize(::oneAPIBackend) = oneAPI.oneL0.synchronize() -KI.supports_float64(::oneAPIBackend) = false # TODO: Check if this is device dependent +KI.supports_float64(::oneAPIBackend) = device_limits().supports_float64 KI.supports_unified(::oneAPIBackend) = true +KI.supports_atomics(::oneAPIBackend) = true KI.functional(::oneAPIBackend) = oneAPI.functional() @@ -41,10 +40,15 @@ end ## Memory Operations 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 @@ -52,15 +56,17 @@ function KI.ndevices(::oneAPIBackend) return length(oneAPI.devices()) end -function KI.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 KI.device!(backend::oneAPIBackend, id::Int) - return oneAPI.device!(id) + oneAPI.device!(id) + return end @@ -68,104 +74,105 @@ end KI.argconvert(::oneAPIBackend, arg) = kernel_convert(arg) -function KI.kernel_function(::oneAPIBackend, f::F, tt::TT=Tuple{}; name = nothing, kwargs...) where {F,TT} - kern = zefunction(f, tt; name, kwargs...) - KI.Kernel{oneAPIBackend, typeof(kern)}(oneAPIBackend(), kern) +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 + zefunction(f, tt; name, backend.always_inline, kwargs...) + end + KI.Kernel{oneAPIBackend, typeof(kern)}(backend, kern) end -function (obj::KI.Kernel{oneAPIBackend})(args...; numworkgroups=(), workgroupsize=(), ndrange=(), max_work_group_size=typemax(Int)) - KI.check_launch_args(numworkgroups, workgroupsize, ndrange) - prod(ndrange) == 0 && return nothing - - numworkgroups, workgroupsize = KI.auto_launch_sizes(obj, numworkgroups, workgroupsize, ndrange, max_work_group_size) - items = (workgroupsize..., ntuple(_ -> 1, 3 - length(workgroupsize))...) - groups = (numworkgroups..., ntuple(_ -> 1, 3 - length(numworkgroups))...) - - obj.kern(args...; items, groups) - return nothing +function KI.launch(obj::KI.Kernel{oneAPIBackend}, groups::Dims{3}, items::Dims{3}, args...; kwargs...) + # 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.kernel_max_work_group_size(kernel::KI.Kernel{<:oneAPIBackend}; max_work_items::Int=typemax(Int))::Int - group_size = oneAPI.launch_configuration(kernel.kern) - Int(min(group_size, max_work_items)) +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 -# querying the device allocates, so cache the limits that every auto-sized launch needs -const DeviceLimits = @NamedTuple{max_work_group_size::Int, max_work_group_dims::NTuple{3, Int}, - max_num_groups::NTuple{3, Int}} -function device_limits() - dev = device()::oneAPI.oneL0.ZeDevice +function KI.launch_configuration(kernel::KI.Kernel{oneAPIBackend}; 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)) + 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 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 -function KI.sub_group_size(::oneAPIBackend)::Int - sg_sizes = oneAPI.oneL0.compute_properties(device()).subGroupSizes - if 32 in sg_sizes - return 32 - elseif 64 in sg_sizes - return 64 - elseif 16 in sg_sizes - return 16 - else - return 1 - end -end +KI.sub_group_size(::oneAPIBackend)::Int = device_limits().sub_group_size function KI.multiprocessor_count(::oneAPIBackend)::Int oneAPI.oneL0.properties(device()).numSlices end -function KI.shfl_down_types(::oneAPIBackend) - res = copy(SPIRVIntrinsics.gentypes) - - res = setdiff(res, [Float64]) - - return res +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 + +# computed with `% T`, which unlike `T(x)` has no error path + @device_override @inline function KI.get_local_id(::Type{T}) where {T} - return (; x = T(get_local_id(1)), y = T(get_local_id(2)), z = T(get_local_id(3))) + return (; x = get_local_id(1) % T, y = get_local_id(2) % T, z = get_local_id(3) % T) end @device_override @inline function KI.get_group_id(::Type{T}) where {T} - return (; x = T(get_group_id(1)), y = T(get_group_id(2)), z = T(get_group_id(3))) -end - -@device_override @inline function KI.get_global_id(::Type{T}) where {T} - return (; x = T(get_global_id(1)), y = T(get_global_id(2)), z = T(get_global_id(3))) + return (; x = get_group_id(1) % T, y = get_group_id(2) % T, z = get_group_id(3) % T) end @device_override @inline function KI.get_local_size(::Type{T}) where {T} - return (; x = T(get_local_size(1)), y = T(get_local_size(2)), z = T(get_local_size(3))) + return (; x = get_local_size(1) % T, y = get_local_size(2) % T, z = get_local_size(3) % T) end @device_override @inline function KI.get_num_groups(::Type{T}) where {T} - return (; x = T(get_num_groups(1)), y = T(get_num_groups(2)), z = T(get_num_groups(3))) -end - -@device_override @inline function KI.get_global_size(::Type{T}) where {T} - return (; x = T(get_global_size(1)), y = T(get_global_size(2)), z = T(get_global_size(3))) + return (; x = get_num_groups(1) % T, y = get_num_groups(2) % T, z = get_num_groups(3) % T) end -@device_override KI.get_sub_group_size() = get_sub_group_size() % UInt32 +@device_override KI.get_sub_group_size(::Type{T}) where {T} = get_sub_group_size() % T -@device_override KI.get_max_sub_group_size() = get_max_sub_group_size() % UInt32 +@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() = get_num_sub_groups() % UInt32 +@device_override KI.get_num_sub_groups(::Type{T}) where {T} = get_num_sub_groups() % T -@device_override KI.get_sub_group_id() = get_sub_group_id() % UInt32 +@device_override KI.get_sub_group_id(::Type{T}) where {T} = get_sub_group_id() % T -@device_override KI.get_sub_group_local_id() = get_sub_group_local_id() % UInt32 +@device_override KI.get_sub_group_local_id(::Type{T}) where {T} = get_sub_group_local_id() % T ## Shared and Scratch Memory diff --git a/test/Project.toml b/test/Project.toml index fe34ca49..e68d74ca 100644 --- a/test/Project.toml +++ b/test/Project.toml @@ -26,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/kernelinterface.jl b/test/kernelinterface.jl index 0854bcfe..75b704c4 100644 --- a/test/kernelinterface.jl +++ b/test/kernelinterface.jl @@ -3,9 +3,4 @@ using oneAPI.oneAPIInterface include(joinpath(dirname(pathof(KernelInterface)), "..", "test", "testsuite.jl")) -skip_tests = Set{String}() -# the events test checks that a waiting task blocks on work queued elsewhere, which is not -# observable when every submission synchronizes (the Aurora LTS workaround) -oneAPI.oneL0.sync_each_submission() && push!(skip_tests, "Events") - -Testsuite.testsuite(oneAPIInterface.oneAPIBackend, "oneAPI", oneAPI, oneArray, oneAPI.oneDeviceArray; skip_tests) +Testsuite.testsuite(oneAPIInterface.oneAPIBackend(), oneArray) From aa98eb171137d39d50a208c73fe1b920f842b210 Mon Sep 17 00:00:00 2001 From: Tim Besard Date: Mon, 28 Sep 2026 10:18:55 +0200 Subject: [PATCH 6/9] Make synchronization cooperative `synchronize` blocked the calling thread in the driver until the work had completed, so no other task could run on it in the meantime, and no other thread could run the GC either. KernelInterface requires `synchronize` to be cooperative. Wait like CUDA.jl does: busy-wait on a non-blocking query first, which keeps the latency of short operations low, and then block in the driver (GC-safe) on one of a few dedicated threads, while the calling task waits for it without blocking the scheduler. This applies to `synchronize()` and `synchronize(::oneStream)`, i.e., to user code, KernelAbstractions, KernelInterface and the synchronizing copies; `synchronize(; blocking=true)` restores the old behavior. The command list and queue methods, which also run from finalizers, keep blocking. --- lib/level-zero/oneL0.jl | 1 + lib/level-zero/synchronization.jl | 195 ++++++++++++++++++++++++++++++ lib/utils/APIUtils.jl | 4 +- src/context.jl | 24 ++-- src/oneAPIKernels.jl | 1 - test/execution.jl | 42 +++++++ 6 files changed, 254 insertions(+), 13 deletions(-) create mode 100644 lib/level-zero/synchronization.jl 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/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/oneAPIKernels.jl b/src/oneAPIKernels.jl index 1ed41b55..47241b72 100644 --- a/src/oneAPIKernels.jl +++ b/src/oneAPIKernels.jl @@ -23,7 +23,6 @@ oneAPIBackend(; prefer_blocks = false, always_inline = false) = oneAPIBackend(pr @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) KI.get_backend(::oneArray) = oneAPIBackend() -# TODO should be non-blocking KI.synchronize(::oneAPIBackend) = oneAPI.oneL0.synchronize() KI.supports_float64(::oneAPIBackend) = device_limits().supports_float64 KI.supports_unified(::oneAPIBackend) = true diff --git a/test/execution.jl b/test/execution.jl index 3efcdbf1..afc02aff 100644 --- a/test/execution.jl +++ b/test/execution.jl @@ -743,6 +743,48 @@ 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() + + # 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. From 760a27a9fede06b1b71e60a6e5dfd909fcc22dff Mon Sep 17 00:00:00 2001 From: Tim Besard Date: Mon, 28 Sep 2026 12:29:42 +0200 Subject: [PATCH 7/9] KernelInterface: specialize KI.launch on its arguments MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Julia doesn't specialize a method on `args...` that it only passes through, which made every launch through KernelInterface's generic launch dispatch dynamically (+1.2 µs and +1 kB per launch on CUDA). --- src/oneAPIKernels.jl | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/src/oneAPIKernels.jl b/src/oneAPIKernels.jl index 47241b72..0f5299c3 100644 --- a/src/oneAPIKernels.jl +++ b/src/oneAPIKernels.jl @@ -84,7 +84,7 @@ function KI.kernel_function(backend::oneAPIBackend, f::F, tt::TT=Tuple{}; name = KI.Kernel{oneAPIBackend, typeof(kern)}(backend, kern) end -function KI.launch(obj::KI.Kernel{oneAPIBackend}, groups::Dims{3}, items::Dims{3}, args...; kwargs...) +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")) From f197c56fe0e72cb95e497bbbb7f79e823f2ba79d Mon Sep 17 00:00:00 2001 From: Tim Besard Date: Wed, 30 Sep 2026 09:44:54 +0200 Subject: [PATCH 8/9] Accept the problem size in KI.launch_configuration KernelInterface 0.3 passes the number of work-items to launch as `nitems`, separately from the bound on the work-group size. --- src/oneAPIKernels.jl | 5 ++++- 1 file changed, 4 insertions(+), 1 deletion(-) diff --git a/src/oneAPIKernels.jl b/src/oneAPIKernels.jl index 0f5299c3..7346df25 100644 --- a/src/oneAPIKernels.jl +++ b/src/oneAPIKernels.jl @@ -98,7 +98,10 @@ function KI.max_work_group_size(kernel::KI.Kernel{oneAPIBackend})::Int # 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}; max_work_group_size::Integer = typemax(Int)) +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 From e93ce2fdb3f6ea6360de06b00ffada96ebab8ca9 Mon Sep 17 00:00:00 2001 From: Tim Besard Date: Wed, 30 Sep 2026 13:39:29 +0200 Subject: [PATCH 9/9] Wait cooperatively after each submission With `ONEAPI_SYNC_EACH_SUBMISSION` set, as on Aurora's LTS stack, every kernel launch blocked the thread in the driver until the kernel had completed. `synchronize` then had nothing left to wait for, so no other task ran in the meantime, which failed the cooperative synchronization test there. Wait for the launch cooperatively, like `synchronize` does. Also warm up the slow path of `synchronize` in that test. Otherwise, in a fresh process, compiling it makes the first calibration step exceed its time budget, and the kernel is then too short to observe other tasks. --- src/compiler/execution.jl | 3 ++- test/execution.jl | 3 +++ 2 files changed, 5 insertions(+), 1 deletion(-) diff --git a/src/compiler/execution.jl b/src/compiler/execution.jl index ed390744..06749f94 100644 --- a/src/compiler/execution.jl +++ b/src/compiler/execution.jl @@ -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/test/execution.jl b/test/execution.jl index afc02aff..330a1612 100644 --- a/test/execution.jl +++ b/test/execution.jl @@ -759,6 +759,9 @@ end 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