From 89b8e446541108e2ac34ed346512f3e539ac0450 Mon Sep 17 00:00:00 2001 From: krasow Date: Mon, 21 Sep 2026 22:48:30 -0500 Subject: [PATCH 1/4] Bloat constraint for halo exchange boundaries --- benchmark/stencil_constraints.jl | 239 ++++++++++++++++++++++ lib/legate_jl_wrapper/include/wrapper.inl | 14 ++ lib/legate_jl_wrapper/src/module.cpp | 7 + lib/legate_jl_wrapper/src/task.cpp | 31 ++- src/api/data.jl | 16 ++ src/api/runtime.jl | 8 +- src/api/tasks.jl | 41 +++- src/ufi.jl | 59 +++++- test/tests/tasking.jl | 64 ++++++ 9 files changed, 458 insertions(+), 21 deletions(-) create mode 100644 benchmark/stencil_constraints.jl diff --git a/benchmark/stencil_constraints.jl b/benchmark/stencil_constraints.jl new file mode 100644 index 0000000..ac678f2 --- /dev/null +++ b/benchmark/stencil_constraints.jl @@ -0,0 +1,239 @@ +using Legate + +const BACKEND = lowercase(get(ENV, "STENCIL_BACKEND", "cpu")) +const GRID_SIZE = parse(Int, get(ENV, "STENCIL_SIZE", "1024")) +const ITERATIONS = parse(Int, get(ENV, "STENCIL_ITERATIONS", "100")) +const WARMUP_ITERATIONS = parse(Int, get(ENV, "STENCIL_WARMUP", "5")) +const SAMPLES = parse(Int, get(ENV, "STENCIL_SAMPLES", "5")) + +if BACKEND == "gpu" + @eval using CUDA +elseif BACKEND != "cpu" + error("STENCIL_BACKEND must be cpu or gpu") +end + +# The field stays fixed so timing isolates repeated partitioning and materialization. +function aligned_stencil(north, south, west, east, center, output) + @inbounds @simd for i in eachindex(output) + output[i] = + 0.2f0 * (north[i] + south[i] + west[i] + east[i] + center[i]) + end + return nothing +end + +function bloated_stencil(core, halo, output) + n = GRID_SIZE + # UFI arrays are locally indexed; monotonic values reveal the clipped halo offset. + delta = round(Int, core[1, 1] - halo[1, 1]) + low_i = delta % n + low_j = delta ÷ n + @inbounds for j in axes(output, 2), i in axes(output, 1) + hi = i + low_i + hj = j + low_j + north = hi > 1 ? halo[hi - 1, hj] : 0.0f0 + south = hi < size(halo, 1) ? halo[hi + 1, hj] : 0.0f0 + west = hj > 1 ? halo[hi, hj - 1] : 0.0f0 + east = hj < size(halo, 2) ? halo[hi, hj + 1] : 0.0f0 + output[i, j] = + 0.2f0 * (north + south + west + east + core[i, j]) + end + return nothing +end + +function aligned_stencil_gpu(args) + north, south, west, east, center, output = args + index = (blockIdx().x - 1) * blockDim().x + threadIdx().x + if index <= length(output) + @inbounds output[index] = + 0.2f0 * + (north[index] + south[index] + west[index] + east[index] + center[index]) + end + return nothing +end + +function bloated_stencil_gpu(args) + core, halo, output = args + index = (blockIdx().x - 1) * blockDim().x + threadIdx().x + if index <= length(output) + n = GRID_SIZE + delta = round(Int, core[1, 1] - halo[1, 1]) + low_i = delta % n + low_j = delta ÷ n + i = (index - 1) % size(output, 1) + 1 + j = (index - 1) ÷ size(output, 1) + 1 + hi = i + low_i + hj = j + low_j + north = hi > 1 ? halo[hi - 1, hj] : 0.0f0 + south = hi < size(halo, 1) ? halo[hi + 1, hj] : 0.0f0 + west = hj > 1 ? halo[hi, hj - 1] : 0.0f0 + east = hj < size(halo, 2) ? halo[hi, hj + 1] : 0.0f0 + @inbounds output[i, j] = + 0.2f0 * (north + south + west + east + core[i, j]) + end + return nothing +end + +function region(array, first, last) + view = Legate.slice(array, 0, first[1], last[1]) + return Legate.slice(view, 1, first[2], last[2]) +end + +function stencil_views(array, n) + center = region(array, (1, 1), (n + 1, n + 1)) + north = region(array, (0, 1), (n, n + 1)) + south = region(array, (2, 1), (n + 2, n + 1)) + west = region(array, (1, 0), (n + 1, n)) + east = region(array, (1, 2), (n + 1, n + 2)) + return (; north, south, west, east, center) +end + +function make_aligned_problem(n) + input_host = zeros(Float32, n + 2, n + 2) + @inbounds for j in 1:n, i in 1:n + input_host[i + 1, j + 1] = Float32(i + n * (j - 1)) + end + + input = Legate.create_array([n + 2, n + 2], Float32) + output = Legate.create_array([n, n], Float32) + copyto!(input, input_host) + views = stencil_views(input, n) + return (; input, output, views) +end + +function make_bloated_problem(n) + input_host = zeros(Float32, n + 2, n + 2) + @inbounds for j in 1:n, i in 1:n + input_host[i + 1, j + 1] = Float32(i + n * (j - 1)) + end + + input = Legate.create_array([n + 2, n + 2], Float32) + output = Legate.create_array([n, n], Float32) + copyto!(input, input_host) + core = region(input, (1, 1), (n + 1, n + 1)) + halo = region(input, (1, 1), (n + 1, n + 1)) + return (; input, core, halo, output) +end + +function submit_aligned!(runtime, library, wrapped, problem) + task = Legate.create_julia_task(runtime, library, wrapped) + inputs = [ + Legate.add_input(task, problem.views.north), + Legate.add_input(task, problem.views.south), + Legate.add_input(task, problem.views.west), + Legate.add_input(task, problem.views.east), + Legate.add_input(task, problem.views.center), + ] + outputs = [Legate.add_output(task, problem.output)] + Legate.default_alignment(task, inputs, outputs) + Legate.submit_task(runtime, task) + return problem +end + +function submit_bloated!(runtime, library, wrapped, problem) + task = Legate.create_julia_task(runtime, library, wrapped) + core = Legate.add_input(task, problem.core) + halo = Legate.add_input(task, problem.halo) + output = Legate.add_output(task, problem.output) + Legate.add_constraint(task, Legate.align(core, output)) + Legate.add_constraint(task, Legate.bloat(core, halo, (1, 1), (1, 1))) + Legate.submit_task(runtime, task) + return problem +end + +function synchronize() + Legate.issue_execution_fence() + Legate.wait_ufi() + return nothing +end + +function run_iterations!(submit!, runtime, library, wrapped, problem, iterations) + for _ in 1:iterations + submit!(runtime, library, wrapped, problem) + end + return problem +end + +function elapsed_sample!(submit!, runtime, library, wrapped, problem, iterations) + synchronize() + start = time_ns() + run_iterations!(submit!, runtime, library, wrapped, problem, iterations) + synchronize() + return (time_ns() - start) / 1.0e9 +end + +function median_value(values) + sorted = sort(values) + middle = length(sorted) ÷ 2 + isodd(length(sorted)) && return sorted[middle + 1] + return (sorted[middle] + sorted[middle + 1]) / 2 +end + +function main() + GRID_SIZE > 0 || error("STENCIL_SIZE must be positive") + ITERATIONS > 0 || error("STENCIL_ITERATIONS must be positive") + SAMPLES > 0 || error("STENCIL_SAMPLES must be positive") + + Legate.Experimental(true) + runtime = Legate.get_runtime() + library = Legate.create_library("stencil_constraint_benchmark") + if BACKEND == "gpu" + aligned = Legate.wrap_task(aligned_stencil_gpu, Legate.GPUBackend) + bloated = Legate.wrap_task(bloated_stencil_gpu, Legate.GPUBackend) + else + aligned = Legate.wrap_task(aligned_stencil, Legate.CPUBackend) + bloated = Legate.wrap_task(bloated_stencil, Legate.CPUBackend) + end + + aligned_problem = make_aligned_problem(GRID_SIZE) + bloated_problem = make_bloated_problem(GRID_SIZE) + + run_iterations!( + submit_aligned!, runtime, library, aligned, aligned_problem, WARMUP_ITERATIONS + ) + run_iterations!( + submit_bloated!, runtime, library, bloated, bloated_problem, WARMUP_ITERATIONS + ) + synchronize() + + aligned_times = Float64[] + bloated_times = Float64[] + for sample in 1:SAMPLES + if isodd(sample) + elapsed = elapsed_sample!( + submit_aligned!, runtime, library, aligned, aligned_problem, ITERATIONS + ) + push!(aligned_times, elapsed) + elapsed = elapsed_sample!( + submit_bloated!, runtime, library, bloated, bloated_problem, ITERATIONS + ) + push!(bloated_times, elapsed) + else + elapsed = elapsed_sample!( + submit_bloated!, runtime, library, bloated, bloated_problem, ITERATIONS + ) + push!(bloated_times, elapsed) + elapsed = elapsed_sample!( + submit_aligned!, runtime, library, aligned, aligned_problem, ITERATIONS + ) + push!(aligned_times, elapsed) + end + end + + aligned_result = Array(aligned_problem.output) + bloated_result = Array(bloated_problem.output) + aligned_result == bloated_result || error("alignment and bloat results differ") + + aligned_median = median_value(aligned_times) + bloated_median = median_value(bloated_times) + println("5-point stencil constraint benchmark") + println(" backend: $BACKEND") + println(" grid: $(GRID_SIZE) x $(GRID_SIZE)") + println(" iterations/sample: $ITERATIONS, samples: $SAMPLES") + println(" default alignment median: $(round(aligned_median; digits=4)) s") + println(" bloat median: $(round(bloated_median; digits=4)) s") + println(" bloat speedup: $(round(aligned_median / bloated_median; digits=3))x") + println(" alignment samples: $aligned_times") + return println(" bloat samples: $bloated_times") +end + +main() diff --git a/lib/legate_jl_wrapper/include/wrapper.inl b/lib/legate_jl_wrapper/include/wrapper.inl index b196cea..cd878d3 100644 --- a/lib/legate_jl_wrapper/include/wrapper.inl +++ b/lib/legate_jl_wrapper/include/wrapper.inl @@ -109,6 +109,20 @@ inline Constraint align(const Variable& a, const Variable& b) { return legate::align(a, b); } +/** + * @ingroup legate_wrapper + * @brief Bloat target partitions around their aligned source partitions. + */ +inline Constraint bloat(const Variable& source, const Variable& target, + const std::vector& low_offsets, + const std::vector& high_offsets) { + return legate::bloat( + source, target, + legate::Span{low_offsets.data(), low_offsets.size()}, + legate::Span{high_offsets.data(), + high_offsets.size()}); +} + /** * @ingroup legate_wrapper * @brief Create an auto task in the runtime. diff --git a/lib/legate_jl_wrapper/src/module.cpp b/lib/legate_jl_wrapper/src/module.cpp index 75cc73c..4566327 100644 --- a/lib/legate_jl_wrapper/src/module.cpp +++ b/lib/legate_jl_wrapper/src/module.cpp @@ -179,6 +179,10 @@ JLCXX_MODULE define_julia_module(jlcxx::Module& mod) { mod.method("slice", [](LogicalStore& s, int32_t dim, legate::Slice sl) { return s.slice(dim, sl); }); + mod.method("slice", + [](LogicalStore& s, int32_t dim, int64_t start, int64_t stop) { + return s.slice(dim, legate::Slice{start, stop}); + }); mod.method( "get_physical_store", [](LogicalStore& s, std::optional target) { @@ -228,6 +232,8 @@ JLCXX_MODULE define_julia_module(jlcxx::Module& mod) { for (int i = 0; i < arr.dim(); i++) result.push_back(s[i]); return result; }); + mod.method("array_from_store", + [](const LogicalStore& store) { return LogicalArray{store}; }); mod.add_type("AutoTask") .method("add_input", static_cast( @@ -294,6 +300,7 @@ JLCXX_MODULE define_julia_module(jlcxx::Module& mod) { mod.method("runtime_sync", &legate_wrapper::runtime::runtime_sync); /* tasking */ mod.method("align", &legate_wrapper::tasking::align); + mod.method("bloat", &legate_wrapper::tasking::bloat); mod.method("domain_from_shape", &legate_wrapper::tasking::domain_from_shape); mod.method("create_manual_task", &legate_wrapper::tasking::create_manual_task); diff --git a/lib/legate_jl_wrapper/src/task.cpp b/lib/legate_jl_wrapper/src/task.cpp index 9bbc304..7550764 100644 --- a/lib/legate_jl_wrapper/src/task.cpp +++ b/lib/legate_jl_wrapper/src/task.cpp @@ -68,16 +68,21 @@ UFI(write, write_accessor); struct ufiFunctor { int* ndim_ptr = nullptr; int64_t* dims_ptr = nullptr; + int64_t* arg_dims_ptr = nullptr; ufiFunctor() = default; - ufiFunctor(int* ndim, int64_t* dims) : ndim_ptr(ndim), dims_ptr(dims) {} + ufiFunctor(int* ndim, int64_t* dims, int64_t* arg_dims) + : ndim_ptr(ndim), dims_ptr(dims), arg_dims_ptr(arg_dims) {} template void operator()(ufi::AccessMode mode, std::uintptr_t& p, int64_t* strides, const legate::PhysicalArray& rf) { + auto shp = rf.shape(); + for (int i = 0; i < DIM && i < REALM_MAX_DIM; ++i) { + arg_dims_ptr[i] = shp.hi[i] - shp.lo[i] + 1; + } if (ndim_ptr && *ndim_ptr == 0) { *ndim_ptr = DIM; - auto shp = rf.shape(); for (int i = 0; i < DIM && i < REALM_MAX_DIM; ++i) { dims_ptr[i] = shp.hi[i] - shp.lo[i] + 1; } @@ -130,6 +135,8 @@ struct UFISlot { char scalar_data[MAX_UFI_ARGS][MAX_SCALAR_SIZE]; int64_t input_strides[MAX_UFI_ARGS][REALM_MAX_DIM]; int64_t output_strides[MAX_UFI_ARGS][REALM_MAX_DIM]; + int64_t input_dims[MAX_UFI_ARGS][REALM_MAX_DIM]; + int64_t output_dims[MAX_UFI_ARGS][REALM_MAX_DIM]; UFISlot() { task_done.store(false); @@ -182,6 +189,16 @@ JULIA_LEGATE_UFI_EXPORT void* legate_get_slot_request_ptr(int slot_id) { return static_cast(&g_ufi_slots[slot_id].request); } +JULIA_LEGATE_UFI_EXPORT int64_t* legate_get_slot_input_dims_ptr(int slot_id) { + if (slot_id < 0 || slot_id >= MAX_UFI_SLOTS) return nullptr; + return &g_ufi_slots[slot_id].input_dims[0][0]; +} + +JULIA_LEGATE_UFI_EXPORT int64_t* legate_get_slot_output_dims_ptr(int slot_id) { + if (slot_id < 0 || slot_id >= MAX_UFI_SLOTS) return nullptr; + return &g_ufi_slots[slot_id].output_dims[0][0]; +} + JULIA_LEGATE_UFI_EXPORT int legate_get_active_call_count() { return g_active_calls.load(); } @@ -273,11 +290,11 @@ inline void JuliaTaskInterface(legate::TaskContext context, bool is_gpu) { slot.request.task_id = task_id; slot.request.ndim = 0; - ufiFunctor functor{&slot.request.ndim, slot.request.dims}; - for (size_t i = 0; i < ni; ++i) { auto ps = context.input(i); std::uintptr_t p = 0; + ufiFunctor functor{no == 0 && i == 0 ? &slot.request.ndim : nullptr, + slot.request.dims, slot.input_dims[i]}; legate::double_dispatch(ps.dim(), ps.type().code(), functor, ufi::AccessMode::READ, p, slot.input_strides[i], ps); @@ -287,6 +304,8 @@ inline void JuliaTaskInterface(legate::TaskContext context, bool is_gpu) { for (size_t i = 0; i < no; ++i) { auto ps = context.output(i); std::uintptr_t p = 0; + ufiFunctor functor{i == 0 ? &slot.request.ndim : nullptr, slot.request.dims, + slot.output_dims[i]}; legate::double_dispatch(ps.dim(), ps.type().code(), functor, ufi::AccessMode::WRITE, p, slot.output_strides[i], ps); @@ -355,6 +374,10 @@ void wrap_ufi(jlcxx::Module& mod) { mod.method("_initialize_async_system", &ufi::initialize_async_system); mod.method("legate_get_max_slots", &ufi::legate_get_max_slots); mod.method("legate_get_slot_request_ptr", &ufi::legate_get_slot_request_ptr); + mod.method("legate_get_slot_input_dims_ptr", + &ufi::legate_get_slot_input_dims_ptr); + mod.method("legate_get_slot_output_dims_ptr", + &ufi::legate_get_slot_output_dims_ptr); mod.method("legate_pop_pending_slot_nonblocking", &ufi::legate_pop_pending_slot_nonblocking); mod.method("legate_get_active_call_count", diff --git a/src/api/data.jl b/src/api/data.jl index f60decb..c0c7653 100644 --- a/src/api/data.jl +++ b/src/api/data.jl @@ -202,6 +202,22 @@ function slice(store::LogicalStore, indices...) return LogicalStore{eltype(store),Int(LegateInternal.dim(impl))}(impl, nothing) end +""" + slice(array::LogicalArray, dim, start, stop) -> LogicalArray + +Return a view of the half-open interval `[start, stop)` in zero-based dimension `dim`. +""" +function slice( + array::LogicalArray{T,N}, dim::Integer, start::Integer, stop::Integer +) where {T,N} + store = LegateInternal.slice( + LegateInternal.data(array.handle), Int32(dim), Int64(start), Int64(stop) + ) + impl = LegateInternal.array_from_store(store) + dims = Tuple(Int.(collect(LegateInternal.shape(impl)))) + return LogicalArray{T,N}(impl, dims, array.order) +end + """ get_physical_store(LogicalStore) -> PhysicalStore get_physical_store(LogicalArray) -> PhysicalStore diff --git a/src/api/runtime.jl b/src/api/runtime.jl index 4496957..83ffaf5 100644 --- a/src/api/runtime.jl +++ b/src/api/runtime.jl @@ -60,7 +60,7 @@ Block until all pending Legate tasks have completed. Useful before reading files written by `h5write` or other async operations. """ -runtime_sync() = LegateInternal.runtime_sync() +runtime_sync() = issue_execution_fence() """ create_library(name::String) -> Library @@ -113,7 +113,11 @@ Issues an execution fence to the runtime. """ function issue_execution_fence(; blocking::Bool=true) if blocking - @threadcall((:legate_issue_execution_fence_blocking, Legate.WRAPPER_LIB_PATH), Cvoid, ()) + @static if VERSION >= v"1.12" + @ccall gc_safe = true WRAPPER_LIB_PATH.legate_issue_execution_fence_blocking()::Cvoid + else + @ccall WRAPPER_LIB_PATH.legate_issue_execution_fence_blocking()::Cvoid + end else LegateInternal.issue_execution_fence(false) end diff --git a/src/api/tasks.jl b/src/api/tasks.jl index 941c164..16862d8 100644 --- a/src/api/tasks.jl +++ b/src/api/tasks.jl @@ -60,6 +60,25 @@ function align(a::Variable, b::Variable) return LegateInternal.align(a, b) end +""" + bloat(source, target, low_offsets, high_offsets) -> Constraint + +Partition `target` like `source`, expanded by the given halo width in each dimension. +""" +function bloat(source::Variable, target::Variable, low_offsets, high_offsets) + isdefined(LegateInternal, :bloat) || + error("bloat constraints require a newer Legate.jl wrapper") + length(low_offsets) == length(high_offsets) || + throw(DimensionMismatch("low and high bloat offsets must have equal lengths")) + all(offset -> offset >= 0, low_offsets) || + throw(ArgumentError("bloat offsets must be nonnegative")) + all(offset -> offset >= 0, high_offsets) || + throw(ArgumentError("bloat offsets must be nonnegative")) + return LegateInternal.bloat( + source, target, to_cxx_vector(low_offsets), to_cxx_vector(high_offsets) + ) +end + function default_alignment( task::LegateTask, inputs::Vector{<:Variable}, outputs::Vector{<:Variable} ) @@ -161,8 +180,16 @@ _gpu_precompile(args...) = nothing function submit_task(rt::CxxPtr{Runtime}, task::LegateTask) drain_pending_frees!() if !isnothing(task.fun) - in_t = Tuple{task.input_types...} - out_t = Tuple{task.output_types...} + n_inputs = length(task.input_types) + in_t = Tuple{ + [Array{T,length(task.arg_dims[i])} for (i, T) in enumerate(task.input_types)]... + } + out_t = Tuple{ + [ + Array{T,length(task.arg_dims[n_inputs + i])} for + (i, T) in enumerate(task.output_types) + ]..., + } sc_t = Tuple{task.scalar_types...} sig = UfiSignature{in_t,out_t,sc_t}() @@ -188,6 +215,8 @@ function submit_task(rt::CxxPtr{Runtime}, task::LegateTask) Ptr{Ptr{Cvoid}}, Ptr{Int64}, Ptr{Int64}, + Ptr{Int64}, + Ptr{Int64}, local_dims_type, typeof(sig), ), @@ -199,11 +228,11 @@ function submit_task(rt::CxxPtr{Runtime}, task::LegateTask) # 2. Precompile the user-provided function with exact types user_arg_types = Any[] - for (T, d) in zip(task.input_types, task.arg_dims) - push!(user_arg_types, Array{T,length(d)}) + for (i, T) in enumerate(task.input_types) + push!(user_arg_types, Array{T,length(task.arg_dims[i])}) end - for (T, d) in zip(task.output_types, task.arg_dims) - push!(user_arg_types, Array{T,length(d)}) + for (i, T) in enumerate(task.output_types) + push!(user_arg_types, Array{T,length(task.arg_dims[n_inputs + i])}) end for T in task.scalar_types push!(user_arg_types, T) diff --git a/src/ufi.jl b/src/ufi.jl index ae3c866..057432b 100644 --- a/src/ufi.jl +++ b/src/ufi.jl @@ -21,6 +21,8 @@ const REALM_MAX_DIM = 6 const MAX_UFI_SLOTS_VAL = 32 const PhysArrPtr = Ptr{Cvoid} const SLOT_REQUEST_PTRS = Vector{Ptr{Cvoid}}(undef, MAX_UFI_SLOTS_VAL) +const SLOT_INPUT_DIMS_PTRS = fill(Ptr{Int64}(C_NULL), MAX_UFI_SLOTS_VAL) +const SLOT_OUTPUT_DIMS_PTRS = fill(Ptr{Int64}(C_NULL), MAX_UFI_SLOTS_VAL) const UFI_INIT_LOCK = ReentrantLock() @@ -76,6 +78,8 @@ struct TaskJob{B<:TaskBackend,M<:UfiMetadata,D<:Tuple} scal_args::Vector{Ptr{Cvoid}} in_strides::Vector{Int64} # flat [arg][REALM_MAX_DIM] element strides out_strides::Vector{Int64} + in_dims::Vector{Int64} # flat [arg][REALM_MAX_DIM] tile dimensions + out_dims::Vector{Int64} local_dims::D meta::M end @@ -161,27 +165,32 @@ end scal_p_ptr::Ptr{Ptr{Cvoid}}, in_str_ptr::Ptr{Int64}, out_str_ptr::Ptr{Int64}, + in_dim_ptr::Ptr{Int64}, + out_dim_ptr::Ptr{Int64}, local_dims::Tuple, ::UfiSignature{InT,OutT,ScT}, ) where {B<:TaskBackend,InT,OutT,ScT} - nd = length(local_dims.parameters) pre = [] # prepare input and output buffers callargs = [] # args passed to the user task post = [] # commit outputs to their destination tiles - # Unroll per-rank stride loads at generation time. - instr(base) = Expr(:tuple, [:(Int(unsafe_load(in_str_ptr, $(base + d)))) for d in 1:nd]...) - outstr(base) = Expr(:tuple, [:(Int(unsafe_load(out_str_ptr, $(base + d)))) for d in 1:nd]...) + metadata(ptr, base, nd) = + Expr(:tuple, [:(Int(unsafe_load($ptr, $(base + d)))) for d in 1:nd]...) for (i, T) in enumerate(InT.parameters) E = eltype(T) + nd = ndims(T) base = (i - 1) * REALM_MAX_DIM sym = gensym(:in) + dsym = gensym(:indims) + ssym = gensym(:instrides) + push!(pre, :($dsym = $(metadata(:in_dim_ptr, base, nd)))) + push!(pre, :($ssym = $(metadata(:in_str_ptr, base, nd)))) push!( pre, :( $sym = _ufi_prepare_input( - backend, $E, unsafe_load(in_p_ptr, $i), local_dims, $(instr(base))) + backend, $E, unsafe_load(in_p_ptr, $i), $dsym, $ssym) ), ) push!(callargs, sym) @@ -189,20 +198,23 @@ end for (i, T) in enumerate(OutT.parameters) E = eltype(T) + nd = ndims(T) base = (i - 1) * REALM_MAX_DIM sym = gensym(:out) state = gensym(:state) + dsym = gensym(:outdims) ssym = gensym(:outs) - push!(pre, :($ssym = $(outstr(base)))) + push!(pre, :($dsym = $(metadata(:out_dim_ptr, base, nd)))) + push!(pre, :($ssym = $(metadata(:out_str_ptr, base, nd)))) push!( pre, :( ($sym, $state) = _ufi_prepare_output( - backend, $E, unsafe_load(out_p_ptr, $i), local_dims, $ssym) + backend, $E, unsafe_load(out_p_ptr, $i), $dsym, $ssym) ), ) push!(callargs, sym) - push!(post, :(_ufi_commit_output(backend, $state, $sym, local_dims, $ssym))) + push!(post, :(_ufi_commit_output(backend, $state, $sym, $dsym, $ssym))) end for (i, T) in enumerate(ScT.parameters) @@ -224,7 +236,9 @@ function _execute_task(job::TaskJob) scalar_args = job.scal_args in_strides = job.in_strides out_strides = job.out_strides - GC.@preserve in_args out_args scalar_args in_strides out_strides begin + in_dims = job.in_dims + out_dims = job.out_dims + GC.@preserve in_args out_args scalar_args in_strides out_strides in_dims out_dims begin _do_call( job.backend, job.meta.fun, @@ -233,6 +247,8 @@ function _execute_task(job::TaskJob) pointer(scalar_args), pointer(in_strides), pointer(out_strides), + pointer(in_dims), + pointer(out_dims), job.local_dims, job.meta.sig, ) @@ -273,6 +289,18 @@ function _copy_strides(ptr::Ptr{Int64}, count::Int) return strides end +function _copy_dims(ptr::Ptr{Int64}, count::Int, common_dims::Tuple) + dims = zeros(Int64, count * REALM_MAX_DIM) + if ptr == C_NULL + for arg in 0:(count - 1), dim in eachindex(common_dims) + dims[arg * REALM_MAX_DIM + dim] = common_dims[dim] + end + elseif count > 0 + unsafe_copyto!(pointer(dims), ptr, length(dims)) + end + return dims +end + # C++ keeps the referenced tile storage valid until the completion callback. function _make_job(slot_id::Int, req::TaskRequest, meta::UfiMetadata) sig_type = typeof(meta.sig) @@ -280,6 +308,8 @@ function _make_job(slot_id::Int, req::TaskRequest, meta::UfiMetadata) out_count = length(sig_type.parameters[2].parameters) scalar_count = length(sig_type.parameters[3].parameters) local_dims = ntuple(i -> Int(max(0, req.dims[i])), Int(req.ndim)) + in_dims_ptr = SLOT_INPUT_DIMS_PTRS[slot_id + 1] + out_dims_ptr = SLOT_OUTPUT_DIMS_PTRS[slot_id + 1] return TaskJob( slot_id, @@ -289,6 +319,8 @@ function _make_job(slot_id::Int, req::TaskRequest, meta::UfiMetadata) _copy_pointer_args(req.scalars_ptr, scalar_count), _copy_strides(req.input_strides_ptr, in_count), _copy_strides(req.output_strides_ptr, out_count), + _copy_dims(in_dims_ptr, in_count, local_dims), + _copy_dims(out_dims_ptr, out_count, local_dims), local_dims, meta, ) @@ -386,6 +418,9 @@ function init_ufi() exit(UFI_ERROR) end + has_arg_dims = + isdefined(LegateInternal, :legate_get_slot_input_dims_ptr) && + isdefined(LegateInternal, :legate_get_slot_output_dims_ptr) for i in 1:max_slots SLOT_REQUEST_PTRS[i] = ccall( (:legate_get_slot_request_ptr, Legate.WRAPPER_LIB_PATH), @@ -393,6 +428,12 @@ function init_ufi() (Cint,), Cint(i-1), ) + if has_arg_dims + SLOT_INPUT_DIMS_PTRS[i] = + LegateInternal.legate_get_slot_input_dims_ptr(Cint(i - 1)).cpp_object + SLOT_OUTPUT_DIMS_PTRS[i] = + LegateInternal.legate_get_slot_output_dims_ptr(Cint(i - 1)).cpp_object + end end LegateInternal._initialize_async_system() diff --git a/test/tests/tasking.jl b/test/tests/tasking.jl index 38ce05f..75b7576 100644 --- a/test/tests/tasking.jl +++ b/test/tests/tasking.jl @@ -58,6 +58,20 @@ function task_scalar(a, b, scalar) end end +function stencil_task(core_indices, halo, output) + offset = Int(core_indices[1] - halo[1]) + @assert length(halo) >= length(output) + @assert 0 <= offset <= 1 + @inbounds for i in eachindex(output) + h = i + offset + center = halo[h] + left = h == 1 ? center : halo[h - 1] + right = h == length(halo) ? center : halo[h + 1] + output[i] = left + 2 * center + right + end + return nothing +end + # get ground truth from base julia base_results = run_base_julia_test() @@ -75,6 +89,7 @@ expected_a = expected_c .* 2.5f0 my_task = Legate.wrap_task(task_test, Legate.CPUBackend) my_4arg_task = Legate.wrap_task(task_4arg, Legate.CPUBackend) my_scalar_task = Legate.wrap_task(task_scalar, Legate.CPUBackend) + my_stencil_task = Legate.wrap_task(stencil_task, Legate.CPUBackend) @testset "Initialization" begin a = Legate.create_array([10, 10], Float32) @@ -94,6 +109,15 @@ expected_a = expected_c .* 2.5f0 @test val_b ≈ base_results.b_init end + @testset "LogicalArray Slice" begin + reference = reshape(Float32.(1:36), 6, 6) + array = Legate.create_array([6, 6], Float32) + copyto!(array, reference) + view = Legate.slice(array, 0, 1, 5) + view = Legate.slice(view, 1, 2, 6) + @test size(view) == (4, 4) + end + a = Legate.create_array([10, 10], Float32) b = Legate.create_array([10, 10], Float32) c = Legate.create_array([10, 10], Float32) @@ -143,4 +167,44 @@ expected_a = expected_c .* 2.5f0 val_a = Array(a) @test val_a ≈ expected_a end + + @testset "Execution Fence" begin + @test Legate.issue_execution_fence() === nothing + @test Legate.runtime_sync() === nothing + end + + @testset "Bloat Constraint (Radius-One Stencil)" begin + if !isdefined(Legate.LegateInternal, :bloat) + @test_skip false + else + n = 100_000 + reference = Float64.(1:n) + core_indices = Legate.create_array([n], Float64) + stencil_input = Legate.create_array([n], Float64) + stencil_output = Legate.create_array([n], Float64) + copyto!(core_indices, reference) + copyto!(stencil_input, reference) + + task = Legate.create_julia_task(rt, lib, my_stencil_task) + core_var = Legate.add_input(task, core_indices) + halo_var = Legate.add_input(task, stencil_input) + output_var = Legate.add_output(task, stencil_output) + Legate.add_constraint(task, Legate.align(core_var, output_var)) + Legate.add_constraint(task, Legate.bloat(core_var, halo_var, (1,), (1,))) + + started_before = Legate.LegateInternal.legate_get_started_count() + Legate.submit_task(rt, task) + result = Array(stencil_output) + started = + Legate.LegateInternal.legate_get_started_count() - started_before + + expected = 4 .* reference + expected[1] = 5 + expected[end] = 4n - 1 + @test result == expected + if Legate.LegateInternal.num_procs() > 1 + @test started > 1 + end + end + end end From 867c9ede6e7df6de05e4b543530112eca0ad88f9 Mon Sep 17 00:00:00 2001 From: krasow Date: Thu, 24 Sep 2026 21:31:16 -0500 Subject: [PATCH 2/4] fix deadlock maybe? --- src/api/tasks.jl | 5 ++--- 1 file changed, 2 insertions(+), 3 deletions(-) diff --git a/src/api/tasks.jl b/src/api/tasks.jl index 16862d8..b61507f 100644 --- a/src/api/tasks.jl +++ b/src/api/tasks.jl @@ -246,14 +246,13 @@ function submit_task(rt::CxxPtr{Runtime}, task::LegateTask) # event loop. Compiling here caches the cubin so the worker's launch is a # cache hit. No-op without CUDA (see CUDAExt). if task.is_gpu - # get_ptr only needs to be GC-safe once a GPU task is in flight; flag it so - # CPU-only programs keep the original accessor path. - LegateInternal.set_gpu_tasking_active(true) _gpu_precompile( task.fun, task.input_types, task.output_types, task.scalar_types, task.arg_dims ) end + # A worker GC would otherwise deadlock against the caller parked in get_ptr. + LegateInternal.set_gpu_tasking_active(true) Threads.atomic_add!(SUBMITTED_COUNT, 1) end return _submit_task(rt, task) From b69e672a5e3229f281bdb6be6912e1ce18b3accb Mon Sep 17 00:00:00 2001 From: krasow Date: Thu, 24 Sep 2026 21:54:49 -0500 Subject: [PATCH 3/4] add broadcast constraint API and tests --- .github/workflows/ci.yml | 2 +- .github/workflows/developer.yml | 2 + benchmark/stencil_constraints.jl | 239 ------------------------------- src/api/tasks.jl | 14 ++ test/tests/tasking.jl | 76 ++++++++++ 5 files changed, 93 insertions(+), 240 deletions(-) delete mode 100644 benchmark/stencil_constraints.jl diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 016b7b1..942a191 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -104,7 +104,7 @@ jobs: LEGATE_SHOW_CONFIG: "1" LEGATE_AUTO_CONFIG: "0" GPUTESTS: "0" # parsed by runtests.jl - LEGATE_CONFIG: "--cpus 1 --utility 1 --sysmem 4000" + LEGATE_CONFIG: "--cpus 2 --utility 1 --sysmem 4000" run: | if julia -e 'exit(VERSION >= v"1.12" ? 0 : 1)'; then export JULIA_NUM_THREADS=1,0 diff --git a/.github/workflows/developer.yml b/.github/workflows/developer.yml index fba4476..4a646d4 100644 --- a/.github/workflows/developer.yml +++ b/.github/workflows/developer.yml @@ -79,6 +79,8 @@ jobs: julia --color=yes -e 'using Pkg; Pkg.build("Legate")' - name: Perform Test + env: + LEGATE_CONFIG: "--cpus 2 --gpus 0 --utility 1 --sysmem 4000" run: | julia --color=yes -e 'using Pkg; Pkg.test("Legate"; julia_args=["--threads=4,1"])' diff --git a/benchmark/stencil_constraints.jl b/benchmark/stencil_constraints.jl deleted file mode 100644 index ac678f2..0000000 --- a/benchmark/stencil_constraints.jl +++ /dev/null @@ -1,239 +0,0 @@ -using Legate - -const BACKEND = lowercase(get(ENV, "STENCIL_BACKEND", "cpu")) -const GRID_SIZE = parse(Int, get(ENV, "STENCIL_SIZE", "1024")) -const ITERATIONS = parse(Int, get(ENV, "STENCIL_ITERATIONS", "100")) -const WARMUP_ITERATIONS = parse(Int, get(ENV, "STENCIL_WARMUP", "5")) -const SAMPLES = parse(Int, get(ENV, "STENCIL_SAMPLES", "5")) - -if BACKEND == "gpu" - @eval using CUDA -elseif BACKEND != "cpu" - error("STENCIL_BACKEND must be cpu or gpu") -end - -# The field stays fixed so timing isolates repeated partitioning and materialization. -function aligned_stencil(north, south, west, east, center, output) - @inbounds @simd for i in eachindex(output) - output[i] = - 0.2f0 * (north[i] + south[i] + west[i] + east[i] + center[i]) - end - return nothing -end - -function bloated_stencil(core, halo, output) - n = GRID_SIZE - # UFI arrays are locally indexed; monotonic values reveal the clipped halo offset. - delta = round(Int, core[1, 1] - halo[1, 1]) - low_i = delta % n - low_j = delta ÷ n - @inbounds for j in axes(output, 2), i in axes(output, 1) - hi = i + low_i - hj = j + low_j - north = hi > 1 ? halo[hi - 1, hj] : 0.0f0 - south = hi < size(halo, 1) ? halo[hi + 1, hj] : 0.0f0 - west = hj > 1 ? halo[hi, hj - 1] : 0.0f0 - east = hj < size(halo, 2) ? halo[hi, hj + 1] : 0.0f0 - output[i, j] = - 0.2f0 * (north + south + west + east + core[i, j]) - end - return nothing -end - -function aligned_stencil_gpu(args) - north, south, west, east, center, output = args - index = (blockIdx().x - 1) * blockDim().x + threadIdx().x - if index <= length(output) - @inbounds output[index] = - 0.2f0 * - (north[index] + south[index] + west[index] + east[index] + center[index]) - end - return nothing -end - -function bloated_stencil_gpu(args) - core, halo, output = args - index = (blockIdx().x - 1) * blockDim().x + threadIdx().x - if index <= length(output) - n = GRID_SIZE - delta = round(Int, core[1, 1] - halo[1, 1]) - low_i = delta % n - low_j = delta ÷ n - i = (index - 1) % size(output, 1) + 1 - j = (index - 1) ÷ size(output, 1) + 1 - hi = i + low_i - hj = j + low_j - north = hi > 1 ? halo[hi - 1, hj] : 0.0f0 - south = hi < size(halo, 1) ? halo[hi + 1, hj] : 0.0f0 - west = hj > 1 ? halo[hi, hj - 1] : 0.0f0 - east = hj < size(halo, 2) ? halo[hi, hj + 1] : 0.0f0 - @inbounds output[i, j] = - 0.2f0 * (north + south + west + east + core[i, j]) - end - return nothing -end - -function region(array, first, last) - view = Legate.slice(array, 0, first[1], last[1]) - return Legate.slice(view, 1, first[2], last[2]) -end - -function stencil_views(array, n) - center = region(array, (1, 1), (n + 1, n + 1)) - north = region(array, (0, 1), (n, n + 1)) - south = region(array, (2, 1), (n + 2, n + 1)) - west = region(array, (1, 0), (n + 1, n)) - east = region(array, (1, 2), (n + 1, n + 2)) - return (; north, south, west, east, center) -end - -function make_aligned_problem(n) - input_host = zeros(Float32, n + 2, n + 2) - @inbounds for j in 1:n, i in 1:n - input_host[i + 1, j + 1] = Float32(i + n * (j - 1)) - end - - input = Legate.create_array([n + 2, n + 2], Float32) - output = Legate.create_array([n, n], Float32) - copyto!(input, input_host) - views = stencil_views(input, n) - return (; input, output, views) -end - -function make_bloated_problem(n) - input_host = zeros(Float32, n + 2, n + 2) - @inbounds for j in 1:n, i in 1:n - input_host[i + 1, j + 1] = Float32(i + n * (j - 1)) - end - - input = Legate.create_array([n + 2, n + 2], Float32) - output = Legate.create_array([n, n], Float32) - copyto!(input, input_host) - core = region(input, (1, 1), (n + 1, n + 1)) - halo = region(input, (1, 1), (n + 1, n + 1)) - return (; input, core, halo, output) -end - -function submit_aligned!(runtime, library, wrapped, problem) - task = Legate.create_julia_task(runtime, library, wrapped) - inputs = [ - Legate.add_input(task, problem.views.north), - Legate.add_input(task, problem.views.south), - Legate.add_input(task, problem.views.west), - Legate.add_input(task, problem.views.east), - Legate.add_input(task, problem.views.center), - ] - outputs = [Legate.add_output(task, problem.output)] - Legate.default_alignment(task, inputs, outputs) - Legate.submit_task(runtime, task) - return problem -end - -function submit_bloated!(runtime, library, wrapped, problem) - task = Legate.create_julia_task(runtime, library, wrapped) - core = Legate.add_input(task, problem.core) - halo = Legate.add_input(task, problem.halo) - output = Legate.add_output(task, problem.output) - Legate.add_constraint(task, Legate.align(core, output)) - Legate.add_constraint(task, Legate.bloat(core, halo, (1, 1), (1, 1))) - Legate.submit_task(runtime, task) - return problem -end - -function synchronize() - Legate.issue_execution_fence() - Legate.wait_ufi() - return nothing -end - -function run_iterations!(submit!, runtime, library, wrapped, problem, iterations) - for _ in 1:iterations - submit!(runtime, library, wrapped, problem) - end - return problem -end - -function elapsed_sample!(submit!, runtime, library, wrapped, problem, iterations) - synchronize() - start = time_ns() - run_iterations!(submit!, runtime, library, wrapped, problem, iterations) - synchronize() - return (time_ns() - start) / 1.0e9 -end - -function median_value(values) - sorted = sort(values) - middle = length(sorted) ÷ 2 - isodd(length(sorted)) && return sorted[middle + 1] - return (sorted[middle] + sorted[middle + 1]) / 2 -end - -function main() - GRID_SIZE > 0 || error("STENCIL_SIZE must be positive") - ITERATIONS > 0 || error("STENCIL_ITERATIONS must be positive") - SAMPLES > 0 || error("STENCIL_SAMPLES must be positive") - - Legate.Experimental(true) - runtime = Legate.get_runtime() - library = Legate.create_library("stencil_constraint_benchmark") - if BACKEND == "gpu" - aligned = Legate.wrap_task(aligned_stencil_gpu, Legate.GPUBackend) - bloated = Legate.wrap_task(bloated_stencil_gpu, Legate.GPUBackend) - else - aligned = Legate.wrap_task(aligned_stencil, Legate.CPUBackend) - bloated = Legate.wrap_task(bloated_stencil, Legate.CPUBackend) - end - - aligned_problem = make_aligned_problem(GRID_SIZE) - bloated_problem = make_bloated_problem(GRID_SIZE) - - run_iterations!( - submit_aligned!, runtime, library, aligned, aligned_problem, WARMUP_ITERATIONS - ) - run_iterations!( - submit_bloated!, runtime, library, bloated, bloated_problem, WARMUP_ITERATIONS - ) - synchronize() - - aligned_times = Float64[] - bloated_times = Float64[] - for sample in 1:SAMPLES - if isodd(sample) - elapsed = elapsed_sample!( - submit_aligned!, runtime, library, aligned, aligned_problem, ITERATIONS - ) - push!(aligned_times, elapsed) - elapsed = elapsed_sample!( - submit_bloated!, runtime, library, bloated, bloated_problem, ITERATIONS - ) - push!(bloated_times, elapsed) - else - elapsed = elapsed_sample!( - submit_bloated!, runtime, library, bloated, bloated_problem, ITERATIONS - ) - push!(bloated_times, elapsed) - elapsed = elapsed_sample!( - submit_aligned!, runtime, library, aligned, aligned_problem, ITERATIONS - ) - push!(aligned_times, elapsed) - end - end - - aligned_result = Array(aligned_problem.output) - bloated_result = Array(bloated_problem.output) - aligned_result == bloated_result || error("alignment and bloat results differ") - - aligned_median = median_value(aligned_times) - bloated_median = median_value(bloated_times) - println("5-point stencil constraint benchmark") - println(" backend: $BACKEND") - println(" grid: $(GRID_SIZE) x $(GRID_SIZE)") - println(" iterations/sample: $ITERATIONS, samples: $SAMPLES") - println(" default alignment median: $(round(aligned_median; digits=4)) s") - println(" bloat median: $(round(bloated_median; digits=4)) s") - println(" bloat speedup: $(round(aligned_median / bloated_median; digits=3))x") - println(" alignment samples: $aligned_times") - return println(" bloat samples: $bloated_times") -end - -main() diff --git a/src/api/tasks.jl b/src/api/tasks.jl index b61507f..b0e7918 100644 --- a/src/api/tasks.jl +++ b/src/api/tasks.jl @@ -79,6 +79,20 @@ function bloat(source::Variable, target::Variable, low_offsets, high_offsets) ) end +""" + broadcast(var) -> Constraint + broadcast(var, axes) -> Constraint + +Give every task the whole of `var`, or only the full extent of zero-based `axes`. +""" +broadcast(var::Variable) = LegateInternal.broadcast(var) + +function broadcast(var::Variable, axes) + all(axis -> axis >= 0, axes) || + throw(ArgumentError("broadcast axes must be nonnegative")) + return LegateInternal.broadcast(var, CxxWrap.StdVector([UInt32(a) for a in axes])) +end + function default_alignment( task::LegateTask, inputs::Vector{<:Variable}, outputs::Vector{<:Variable} ) diff --git a/test/tests/tasking.jl b/test/tests/tasking.jl index 75b7576..cd6adb8 100644 --- a/test/tests/tasking.jl +++ b/test/tests/tasking.jl @@ -72,6 +72,26 @@ function stencil_task(core_indices, halo, output) return nothing end +function gather_task(indices, table, out) + @inbounds for i in eachindex(out) + out[i] = table[Int(indices[i])] + end + return nothing +end + +function row_normalize_task(a, out) + @inbounds for i in axes(a, 1) + s = zero(eltype(a)) + for j in axes(a, 2) + s += a[i, j] + end + for j in axes(a, 2) + out[i, j] = a[i, j] / s + end + end + return nothing +end + # get ground truth from base julia base_results = run_base_julia_test() @@ -90,6 +110,8 @@ expected_a = expected_c .* 2.5f0 my_4arg_task = Legate.wrap_task(task_4arg, Legate.CPUBackend) my_scalar_task = Legate.wrap_task(task_scalar, Legate.CPUBackend) my_stencil_task = Legate.wrap_task(stencil_task, Legate.CPUBackend) + my_gather_task = Legate.wrap_task(gather_task, Legate.CPUBackend) + my_row_normalize_task = Legate.wrap_task(row_normalize_task, Legate.CPUBackend) @testset "Initialization" begin a = Legate.create_array([10, 10], Float32) @@ -207,4 +229,58 @@ expected_a = expected_c .* 2.5f0 end end end + + # Operands must be large enough that Legate would otherwise split them. + @testset "Broadcast Constraint (Full)" begin + n = 100_000 + indices_host = Float64.(n:-1:1) + table_host = Float64.(1:n) .^ 2 + indices = Legate.create_array([n], Float64) + table = Legate.create_array([n], Float64) + out = Legate.create_array([n], Float64) + copyto!(indices, indices_host) + copyto!(table, table_host) + + task = Legate.create_julia_task(rt, lib, my_gather_task) + indices_var = Legate.add_input(task, indices) + table_var = Legate.add_input(task, table) + out_var = Legate.add_output(task, out) + Legate.add_constraint(task, Legate.align(indices_var, out_var)) + Legate.add_constraint(task, Legate.broadcast(table_var)) + + started_before = Legate.LegateInternal.legate_get_started_count() + Legate.submit_task(rt, task) + result = Array(out) + started = Legate.LegateInternal.legate_get_started_count() - started_before + + @test result == reverse(table_host) + if Legate.LegateInternal.num_procs() > 1 + @test started > 1 + end + end + + @testset "Broadcast Constraint (Axes)" begin + m, k = 1024, 1024 + a_host = reshape(Float64.(1:(m * k)), m, k) + a = Legate.create_array([m, k], Float64) + out = Legate.create_array([m, k], Float64) + copyto!(a, a_host) + + task = Legate.create_julia_task(rt, lib, my_row_normalize_task) + a_var = Legate.add_input(task, a) + out_var = Legate.add_output(task, out) + Legate.add_constraint(task, Legate.align(a_var, out_var)) + # Row sums need every column in each tile. + Legate.add_constraint(task, Legate.broadcast(a_var, (1,))) + + started_before = Legate.LegateInternal.legate_get_started_count() + Legate.submit_task(rt, task) + result = Array(out) + started = Legate.LegateInternal.legate_get_started_count() - started_before + + @test result ≈ a_host ./ sum(a_host; dims=2) + if Legate.LegateInternal.num_procs() > 1 + @test started > 1 + end + end end From ccb92606ed18a44f0e1481c0b2bf29c4e8a69471 Mon Sep 17 00:00:00 2001 From: krasow Date: Thu, 24 Sep 2026 22:08:50 -0500 Subject: [PATCH 4/4] test: reorg tasking testing infra --- test/runtests.jl | 6 +- test/tests/{tasking.jl => tasking/basic.jl} | 129 ----------------- test/tests/tasking/constraints.jl | 130 ++++++++++++++++++ test/tests/{tasking_gpu.jl => tasking/gpu.jl} | 0 test/tests/tasking/setup.jl | 3 + 5 files changed, 137 insertions(+), 131 deletions(-) rename test/tests/{tasking.jl => tasking/basic.jl} (51%) create mode 100644 test/tests/tasking/constraints.jl rename test/tests/{tasking_gpu.jl => tasking/gpu.jl} (100%) create mode 100644 test/tests/tasking/setup.jl diff --git a/test/runtests.jl b/test/runtests.jl index 11e5ccb..812718e 100644 --- a/test/runtests.jl +++ b/test/runtests.jl @@ -61,7 +61,9 @@ include("tests/hdf5.jl") include("tests/stability.jl") include("tests/basic.jl") -include("tests/tasking.jl") +include("tests/tasking/setup.jl") +include("tests/tasking/basic.jl") +include("tests/tasking/constraints.jl") if run_gpu_tests - include("tests/tasking_gpu.jl") + include("tests/tasking/gpu.jl") end diff --git a/test/tests/tasking.jl b/test/tests/tasking/basic.jl similarity index 51% rename from test/tests/tasking.jl rename to test/tests/tasking/basic.jl index cd6adb8..9ec5d50 100644 --- a/test/tests/tasking.jl +++ b/test/tests/tasking/basic.jl @@ -58,40 +58,6 @@ function task_scalar(a, b, scalar) end end -function stencil_task(core_indices, halo, output) - offset = Int(core_indices[1] - halo[1]) - @assert length(halo) >= length(output) - @assert 0 <= offset <= 1 - @inbounds for i in eachindex(output) - h = i + offset - center = halo[h] - left = h == 1 ? center : halo[h - 1] - right = h == length(halo) ? center : halo[h + 1] - output[i] = left + 2 * center + right - end - return nothing -end - -function gather_task(indices, table, out) - @inbounds for i in eachindex(out) - out[i] = table[Int(indices[i])] - end - return nothing -end - -function row_normalize_task(a, out) - @inbounds for i in axes(a, 1) - s = zero(eltype(a)) - for j in axes(a, 2) - s += a[i, j] - end - for j in axes(a, 2) - out[i, j] = a[i, j] / s - end - end - return nothing -end - # get ground truth from base julia base_results = run_base_julia_test() @@ -103,15 +69,9 @@ expected_a = expected_c .* 2.5f0 @testset verbose=true "CPU Tasking" begin Legate.Experimental(true) # tasking is experimental - rt = Legate.get_runtime() - lib = Legate.create_library("test_comparison") - my_task = Legate.wrap_task(task_test, Legate.CPUBackend) my_4arg_task = Legate.wrap_task(task_4arg, Legate.CPUBackend) my_scalar_task = Legate.wrap_task(task_scalar, Legate.CPUBackend) - my_stencil_task = Legate.wrap_task(stencil_task, Legate.CPUBackend) - my_gather_task = Legate.wrap_task(gather_task, Legate.CPUBackend) - my_row_normalize_task = Legate.wrap_task(row_normalize_task, Legate.CPUBackend) @testset "Initialization" begin a = Legate.create_array([10, 10], Float32) @@ -194,93 +154,4 @@ expected_a = expected_c .* 2.5f0 @test Legate.issue_execution_fence() === nothing @test Legate.runtime_sync() === nothing end - - @testset "Bloat Constraint (Radius-One Stencil)" begin - if !isdefined(Legate.LegateInternal, :bloat) - @test_skip false - else - n = 100_000 - reference = Float64.(1:n) - core_indices = Legate.create_array([n], Float64) - stencil_input = Legate.create_array([n], Float64) - stencil_output = Legate.create_array([n], Float64) - copyto!(core_indices, reference) - copyto!(stencil_input, reference) - - task = Legate.create_julia_task(rt, lib, my_stencil_task) - core_var = Legate.add_input(task, core_indices) - halo_var = Legate.add_input(task, stencil_input) - output_var = Legate.add_output(task, stencil_output) - Legate.add_constraint(task, Legate.align(core_var, output_var)) - Legate.add_constraint(task, Legate.bloat(core_var, halo_var, (1,), (1,))) - - started_before = Legate.LegateInternal.legate_get_started_count() - Legate.submit_task(rt, task) - result = Array(stencil_output) - started = - Legate.LegateInternal.legate_get_started_count() - started_before - - expected = 4 .* reference - expected[1] = 5 - expected[end] = 4n - 1 - @test result == expected - if Legate.LegateInternal.num_procs() > 1 - @test started > 1 - end - end - end - - # Operands must be large enough that Legate would otherwise split them. - @testset "Broadcast Constraint (Full)" begin - n = 100_000 - indices_host = Float64.(n:-1:1) - table_host = Float64.(1:n) .^ 2 - indices = Legate.create_array([n], Float64) - table = Legate.create_array([n], Float64) - out = Legate.create_array([n], Float64) - copyto!(indices, indices_host) - copyto!(table, table_host) - - task = Legate.create_julia_task(rt, lib, my_gather_task) - indices_var = Legate.add_input(task, indices) - table_var = Legate.add_input(task, table) - out_var = Legate.add_output(task, out) - Legate.add_constraint(task, Legate.align(indices_var, out_var)) - Legate.add_constraint(task, Legate.broadcast(table_var)) - - started_before = Legate.LegateInternal.legate_get_started_count() - Legate.submit_task(rt, task) - result = Array(out) - started = Legate.LegateInternal.legate_get_started_count() - started_before - - @test result == reverse(table_host) - if Legate.LegateInternal.num_procs() > 1 - @test started > 1 - end - end - - @testset "Broadcast Constraint (Axes)" begin - m, k = 1024, 1024 - a_host = reshape(Float64.(1:(m * k)), m, k) - a = Legate.create_array([m, k], Float64) - out = Legate.create_array([m, k], Float64) - copyto!(a, a_host) - - task = Legate.create_julia_task(rt, lib, my_row_normalize_task) - a_var = Legate.add_input(task, a) - out_var = Legate.add_output(task, out) - Legate.add_constraint(task, Legate.align(a_var, out_var)) - # Row sums need every column in each tile. - Legate.add_constraint(task, Legate.broadcast(a_var, (1,))) - - started_before = Legate.LegateInternal.legate_get_started_count() - Legate.submit_task(rt, task) - result = Array(out) - started = Legate.LegateInternal.legate_get_started_count() - started_before - - @test result ≈ a_host ./ sum(a_host; dims=2) - if Legate.LegateInternal.num_procs() > 1 - @test started > 1 - end - end end diff --git a/test/tests/tasking/constraints.jl b/test/tests/tasking/constraints.jl new file mode 100644 index 0000000..e875b61 --- /dev/null +++ b/test/tests/tasking/constraints.jl @@ -0,0 +1,130 @@ +# Operands must be large enough that Legate would otherwise split them. + +function stencil_task(core_indices, halo, output) + offset = Int(core_indices[1] - halo[1]) + @assert length(halo) >= length(output) + @assert 0 <= offset <= 1 + @inbounds for i in eachindex(output) + h = i + offset + center = halo[h] + left = h == 1 ? center : halo[h - 1] + right = h == length(halo) ? center : halo[h + 1] + output[i] = left + 2 * center + right + end + return nothing +end + +function gather_task(indices, table, out) + @inbounds for i in eachindex(out) + out[i] = table[Int(indices[i])] + end + return nothing +end + +function row_normalize_task(a, out) + @inbounds for i in axes(a, 1) + s = zero(eltype(a)) + for j in axes(a, 2) + s += a[i, j] + end + for j in axes(a, 2) + out[i, j] = a[i, j] / s + end + end + return nothing +end + +@testset verbose=true "Tasking Constraints" begin + Legate.Experimental(true) # tasking is experimental + my_stencil_task = Legate.wrap_task(stencil_task, Legate.CPUBackend) + my_gather_task = Legate.wrap_task(gather_task, Legate.CPUBackend) + my_row_normalize_task = Legate.wrap_task(row_normalize_task, Legate.CPUBackend) + + @testset "Bloat Constraint (Radius-One Stencil)" begin + if !isdefined(Legate.LegateInternal, :bloat) + @test_skip false + else + n = 100_000 + reference = Float64.(1:n) + core_indices = Legate.create_array([n], Float64) + stencil_input = Legate.create_array([n], Float64) + stencil_output = Legate.create_array([n], Float64) + copyto!(core_indices, reference) + copyto!(stencil_input, reference) + + task = Legate.create_julia_task(rt, lib, my_stencil_task) + core_var = Legate.add_input(task, core_indices) + halo_var = Legate.add_input(task, stencil_input) + output_var = Legate.add_output(task, stencil_output) + Legate.add_constraint(task, Legate.align(core_var, output_var)) + Legate.add_constraint(task, Legate.bloat(core_var, halo_var, (1,), (1,))) + + started_before = Legate.LegateInternal.legate_get_started_count() + Legate.submit_task(rt, task) + result = Array(stencil_output) + started = + Legate.LegateInternal.legate_get_started_count() - started_before + + expected = 4 .* reference + expected[1] = 5 + expected[end] = 4n - 1 + @test result == expected + if Legate.LegateInternal.num_procs() > 1 + @test started > 1 + end + end + end + + @testset "Broadcast Constraint (Full)" begin + n = 100_000 + indices_host = Float64.(n:-1:1) + table_host = Float64.(1:n) .^ 2 + indices = Legate.create_array([n], Float64) + table = Legate.create_array([n], Float64) + out = Legate.create_array([n], Float64) + copyto!(indices, indices_host) + copyto!(table, table_host) + + task = Legate.create_julia_task(rt, lib, my_gather_task) + indices_var = Legate.add_input(task, indices) + table_var = Legate.add_input(task, table) + out_var = Legate.add_output(task, out) + Legate.add_constraint(task, Legate.align(indices_var, out_var)) + Legate.add_constraint(task, Legate.broadcast(table_var)) + + started_before = Legate.LegateInternal.legate_get_started_count() + Legate.submit_task(rt, task) + result = Array(out) + started = Legate.LegateInternal.legate_get_started_count() - started_before + + @test result == reverse(table_host) + if Legate.LegateInternal.num_procs() > 1 + @test started > 1 + end + end + + @testset "Broadcast Constraint (Axes)" begin + m, k = 1024, 1024 + a_host = reshape(Float64.(1:(m * k)), m, k) + a = Legate.create_array([m, k], Float64) + out = Legate.create_array([m, k], Float64) + copyto!(a, a_host) + + task = Legate.create_julia_task(rt, lib, my_row_normalize_task) + a_var = Legate.add_input(task, a) + out_var = Legate.add_output(task, out) + Legate.add_constraint(task, Legate.align(a_var, out_var)) + # Row sums need every column in each tile. + Legate.add_constraint(task, Legate.broadcast(a_var, (1,))) + + started_before = Legate.LegateInternal.legate_get_started_count() + Legate.submit_task(rt, task) + result = Array(out) + started = Legate.LegateInternal.legate_get_started_count() - started_before + + @test result ≈ a_host ./ sum(a_host; dims=2) + if Legate.LegateInternal.num_procs() > 1 + @test started > 1 + end + end +end diff --git a/test/tests/tasking_gpu.jl b/test/tests/tasking/gpu.jl similarity index 100% rename from test/tests/tasking_gpu.jl rename to test/tests/tasking/gpu.jl diff --git a/test/tests/tasking/setup.jl b/test/tests/tasking/setup.jl new file mode 100644 index 0000000..5bbc079 --- /dev/null +++ b/test/tests/tasking/setup.jl @@ -0,0 +1,3 @@ +# Shared by the tasking test files. +const rt = Legate.get_runtime() +const lib = Legate.create_library("test_tasking")