Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion .github/workflows/ci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
2 changes: 2 additions & 0 deletions .github/workflows/developer.yml
Original file line number Diff line number Diff line change
Expand Up @@ -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"])'

Expand Down
14 changes: 14 additions & 0 deletions lib/legate_jl_wrapper/include/wrapper.inl
Original file line number Diff line number Diff line change
Expand Up @@ -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<std::uint64_t>& low_offsets,
const std::vector<std::uint64_t>& high_offsets) {
return legate::bloat(
source, target,
legate::Span<const std::uint64_t>{low_offsets.data(), low_offsets.size()},
legate::Span<const std::uint64_t>{high_offsets.data(),
high_offsets.size()});
}

/**
* @ingroup legate_wrapper
* @brief Create an auto task in the runtime.
Expand Down
7 changes: 7 additions & 0 deletions lib/legate_jl_wrapper/src/module.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -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<legate::mapping::StoreTarget> target) {
Expand Down Expand Up @@ -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>("AutoTask")
.method("add_input", static_cast<Variable (AutoTask::*)(LogicalArray)>(
Expand Down Expand Up @@ -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);
Expand Down
31 changes: 27 additions & 4 deletions lib/legate_jl_wrapper/src/task.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -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 <legate::Type::Code CODE, int DIM>
void operator()(ufi::AccessMode mode, std::uintptr_t& p, int64_t* strides,
const legate::PhysicalArray& rf) {
auto shp = rf.shape<DIM>();
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<DIM>();
for (int i = 0; i < DIM && i < REALM_MAX_DIM; ++i) {
dims_ptr[i] = shp.hi[i] - shp.lo[i] + 1;
}
Expand Down Expand Up @@ -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);
Expand Down Expand Up @@ -182,6 +189,16 @@ JULIA_LEGATE_UFI_EXPORT void* legate_get_slot_request_ptr(int slot_id) {
return static_cast<void*>(&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();
}
Expand Down Expand Up @@ -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);
Expand All @@ -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);
Expand Down Expand Up @@ -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",
Expand Down
16 changes: 16 additions & 0 deletions src/api/data.jl
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
8 changes: 6 additions & 2 deletions src/api/runtime.jl
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down
60 changes: 51 additions & 9 deletions src/api/tasks.jl
Original file line number Diff line number Diff line change
Expand Up @@ -60,6 +60,39 @@ 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

"""
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}
)
Expand Down Expand Up @@ -161,8 +194,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}()
Expand All @@ -188,6 +229,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),
),
Expand All @@ -199,11 +242,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)
Expand All @@ -217,14 +260,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)
Expand Down
Loading
Loading