backend/native API
The package Luna-Flow/luna_thread/backend/native, imported as @native,
is the C foreign function interface backend. It runs the integer kernels of the
C runtime on MoonBit arrays, submits workflows to the C scheduler, and turns
plans into typed native requests. It builds only for the native target and
links the C runtime from its native stubs.
The enumerations of this package are read-only outside it: match on their
constructors. The request records are built by the *_from_plan functions.
Backend information
backend_target
Returns Native.
pub fn backend_target() -> @shared.BackendTarget
default_plan_kind
Returns Map.
pub fn default_plan_kind() -> @plan.PlanKind
package_name
Returns "backend/native".
pub fn package_name() -> String
supports_openmp
Returns true.
pub fn supports_openmp() -> Bool
The value is a constant. The moon build compiles the C stubs without OpenMP,
so in that build the kernels run on the calling thread; see the
native backend design.
Kernels
All kernels take a FixedArray[Int] of length , a worker count and a
chunk size , block until they finish, and leave the input unchanged. They
require
and fail with InvalidArgument otherwise. The facade functions of the same
names in the core API wrap these with optional arguments.
execute_map_i32
Returns a new array with every element doubled, or Overflow when some
does not fit in 32 bits.
pub fn execute_map_i32(FixedArray[Int], Int, Int) -> Result[FixedArray[Int], NativeRequestError]
execute_reduce_sum_i32
Returns the sum of the elements, or 0 on any failure.
pub fn execute_reduce_sum_i32(FixedArray[Int], Int, Int) -> Int
The kernel adds each chunk from the left and then adds the chunk sums from the
left, failing when any of these partial sums overflows. A returned non-zero
value is the exact sum; 0 is either the sum or a failure.
execute_scan_sum_i32
Returns the inclusive prefix sums.
pub fn execute_scan_sum_i32(FixedArray[Int], Int, Int) -> Result[FixedArray[Int], NativeRequestError]
It fails with Overflow when a partial sum overflows.
test "kernels" {
let input : FixedArray[Int] = [3, 1, 4, 1, 5, 9]
debug_inspect(
@native.execute_scan_sum_i32(input, 3, 2),
content="Ok(<FixedArray: [3, 4, 8, 9, 14, 23]>)",
)
assert_eq(@native.execute_reduce_sum_i32(input, 3, 2), 23)
debug_inspect(@native.execute_map_i32(input, 2, 2), content="Err(InvalidArgument)")
}
Workflows
workflow_is_supported
Returns true when @workflow.validate finds no issue and the workflow’s
policy names the native backend.
pub fn workflow_is_supported(@workflow.Workflow) -> Bool
submit_workflow
Returns @workflow.submit(workflow, backend=Native): a validated submission
record. It does not run the workflow.
pub fn submit_workflow(@workflow.Workflow) -> @workflow.Submission
WorkflowRuntimeState
The state of a workflow in the C runtime.
pub enum WorkflowRuntimeState {
Submitted
Running
Completed
Failed
Rejected
} derive(Eq, @debug.Debug)
Submitted holds until a worker takes the first node. The runtime never sets
Rejected; a rejected submission yields an empty handle instead.
WorkflowResult
A snapshot of a workflow: its state, the runtime status code, the number of
completed nodes and the id of the failed node (-1 when none failed).
pub struct WorkflowResult {
state : WorkflowRuntimeState
status : Int
completed_nodes : Int
failed_node_id : Int
} derive(Eq, @debug.Debug)
status is one of the codes in the architecture guide.
NativeWorkflowHandle and WorkflowHandle
An opaque pointer to a running workflow in the C runtime, and its MoonBit wrapper.
#external
pub type NativeWorkflowHandle
pub struct WorkflowHandle {
raw : NativeWorkflowHandle
}
submit_workflow_async
Copies a workflow into flat integer arrays and starts it on the C scheduler
with worker_count threads from its policy.
pub fn submit_workflow_async(@workflow.Workflow) -> Result[WorkflowHandle, NativeRequestError]
The workflow is not validated in MoonBit; the C runtime checks it and, when it
rejects the graph, returns an empty handle. The function always returns Ok.
Compute nodes are scheduled, but their plans are not executed.
poll_workflow
Returns a snapshot without waiting.
pub fn poll_workflow(WorkflowHandle) -> WorkflowResult
wait_workflow
Waits until the workflow is Completed or Failed and returns the snapshot.
pub fn wait_workflow(WorkflowHandle) -> WorkflowResult
On an empty handle it returns at once with state Submitted and status 7.
It does not return if the workflow deadlocks.
drop_workflow
Stops and joins the runtime threads and frees the workflow. Call it once per handle; it does nothing for an empty handle.
pub fn drop_workflow(WorkflowHandle) -> Unit
Typed native requests
NativeValueType and NativeReductionKernel
The value types and reduction kernels the C runtime implements.
pub enum NativeValueType {
I32
I64
} derive(Eq, @debug.Debug)
pub enum NativeReductionKernel {
Sum
Min
Max
} derive(Eq, @debug.Debug)
NativeRequestError
Why a plan could not become a native request, or why a kernel failed.
pub enum NativeRequestError {
UnsupportedBackend
UnsupportedMode
UnsupportedValueType
UnsupportedReductionKernel
InvalidArgument
Overflow
} derive(Eq, @debug.Debug)
NativeMapRequest, NativeReduceRequest and NativeScanRequest
Typed counterparts of the C request structs.
pub struct NativeMapRequest {
input : @shared.NativeBuffer
output : @shared.NativeBuffer
element_count : Int
value_type : NativeValueType
worker_count : Int
chunk_size : Int
} derive(Eq, @debug.Debug)
pub struct NativeReduceRequest {
input : @shared.NativeBuffer
output : @shared.NativeBuffer
element_count : Int
value_type : NativeValueType
reduction_kernel : NativeReductionKernel
worker_count : Int
chunk_size : Int
} derive(Eq, @debug.Debug)
pub struct NativeScanRequest {
input : @shared.NativeBuffer
output : @shared.NativeBuffer
element_count : Int
value_type : NativeValueType
reduction_kernel : NativeReductionKernel
worker_count : Int
chunk_size : Int
} derive(Eq, @debug.Debug)
The input and output buffers of requests built from plans are
NativeBuffer::new(0, 0): a plan describes the shape of the data, not the
data. No function of the module sends these records to C; they are the
validated description that a future executor will fill with buffers.
map_request_from_plan
Turns a Map plan into a NativeMapRequest.
pub fn map_request_from_plan(@plan.Plan) -> Result[NativeMapRequest, NativeRequestError]
It fails with UnsupportedMode when the plan is not a Map plan, and
otherwise with the error mapped from the first issue of @plan.validate:
value type, reduction kernel, backend and mode issues map to the error of the
same name, MissingReductionKernel maps to UnsupportedReductionKernel, and
the remaining issues map to InvalidArgument.
reduce_request_from_plan
Turns a Reduce or MapReduce plan into a NativeReduceRequest, with the same
error rules.
pub fn reduce_request_from_plan(@plan.Plan) -> Result[NativeReduceRequest, NativeRequestError]
scan_request_from_plan
Turns a Scan plan into a NativeScanRequest whose kernel is Sum, with the
same error rules.
pub fn scan_request_from_plan(@plan.Plan) -> Result[NativeScanRequest, NativeRequestError]
map_request_is_valid, reduce_request_is_valid and scan_request_is_valid
Check the request fields the C validator also checks: positive element count,
worker count and chunk size, a supported value type, and for scans the Sum
kernel.
pub fn map_request_is_valid(NativeMapRequest) -> Bool
pub fn reduce_request_is_valid(NativeReduceRequest) -> Bool
pub fn scan_request_is_valid(NativeScanRequest) -> Bool
map_request_value_type, reduce_request_value_type, scan_request_value_type, reduce_request_kernel and scan_request_kernel
Return the C codes of a request’s value type (0 for I32, 1 for I64)
and kernel (0 for Sum, 1 for Min, 2 for Max).
pub fn map_request_value_type(NativeMapRequest) -> Int
pub fn reduce_request_value_type(NativeReduceRequest) -> Int
pub fn scan_request_value_type(NativeScanRequest) -> Int
pub fn reduce_request_kernel(NativeReduceRequest) -> Int
pub fn scan_request_kernel(NativeScanRequest) -> Int
test "requests from plans" {
let policy = @shared.make_execution_policy(worker_count=2, chunk_size=4).unwrap()
let plan = @plan.map_reduce("sum", @plan.i64_type(), 8, @plan.max_reduction(), policy~)
guard @native.reduce_request_from_plan(plan) is Ok(request) else {
fail("expected a request")
}
assert_true(@native.reduce_request_is_valid(request))
assert_eq(@native.reduce_request_value_type(request), 1)
assert_eq(@native.reduce_request_kernel(request), 2)
let scan = @plan.scan("prefix", @plan.i32_type(), 8, policy~)
debug_inspect(@native.map_request_from_plan(scan), content="Err(UnsupportedMode)")
}
Raw foreign functions
These extern "C" declarations are public so that callers can reach kernels
the typed functions do not wrap, such as the Int64 kernels and minimum and
maximum reductions. Arrays are borrowed for the duration of the call. Map and
scan return a status code and write into output, which must have at least
length elements; reductions return the result, or 0 on failure. The
arguments are the input, its length, the worker count and the chunk size.
ffi_execute_map_i32 and ffi_execute_map_i64
Double every element into output.
pub fn ffi_execute_map_i32(FixedArray[Int], Int, Int, Int, FixedArray[Int]) -> Int
pub fn ffi_execute_map_i64(FixedArray[Int64], Int, Int, Int, FixedArray[Int64]) -> Int
ffi_execute_reduce_sum_i32, ffi_execute_reduce_sum_i64, ffi_execute_reduce_min_i32, ffi_execute_reduce_min_i64, ffi_execute_reduce_max_i32 and ffi_execute_reduce_max_i64
Return the sum, minimum or maximum of the input.
pub fn ffi_execute_reduce_sum_i32(FixedArray[Int], Int, Int, Int) -> Int
pub fn ffi_execute_reduce_sum_i64(FixedArray[Int64], Int, Int, Int) -> Int64
pub fn ffi_execute_reduce_min_i32(FixedArray[Int], Int, Int, Int) -> Int
pub fn ffi_execute_reduce_min_i64(FixedArray[Int64], Int, Int, Int) -> Int64
pub fn ffi_execute_reduce_max_i32(FixedArray[Int], Int, Int, Int) -> Int
pub fn ffi_execute_reduce_max_i64(FixedArray[Int64], Int, Int, Int) -> Int64
Minimum and maximum never overflow.
ffi_execute_scan_sum_i32 and ffi_execute_scan_sum_i64
Write the inclusive prefix sums into output.
pub fn ffi_execute_scan_sum_i32(FixedArray[Int], Int, Int, Int, FixedArray[Int]) -> Int
pub fn ffi_execute_scan_sum_i64(FixedArray[Int64], Int, Int, Int, FixedArray[Int64]) -> Int
ffi_submit_workflow_async
Starts a workflow from flat arrays: worker count; capability ids, kind codes
and count; node ids, kind codes, capability ids (-1 for none) and count; edge
sources, targets, kind codes and count. Returns a null handle when the runtime
rejects the graph.
pub fn ffi_submit_workflow_async(Int, FixedArray[Int], FixedArray[Int], Int, FixedArray[Int], FixedArray[Int], FixedArray[Int], Int, FixedArray[Int], FixedArray[Int], FixedArray[Int], Int) -> NativeWorkflowHandle
ffi_workflow_poll_state, ffi_workflow_poll_status, ffi_workflow_poll_completed_nodes, ffi_workflow_poll_failed_node_id, ffi_workflow_wait_status and ffi_workflow_destroy
Read one field of a workflow snapshot, wait and return the final status, or destroy the workflow.
pub fn ffi_workflow_poll_state(NativeWorkflowHandle) -> Int
pub fn ffi_workflow_poll_status(NativeWorkflowHandle) -> Int
pub fn ffi_workflow_poll_completed_nodes(NativeWorkflowHandle) -> Int
pub fn ffi_workflow_poll_failed_node_id(NativeWorkflowHandle) -> Int
pub fn ffi_workflow_wait_status(NativeWorkflowHandle) -> Int
pub fn ffi_workflow_destroy(NativeWorkflowHandle) -> Unit
On a null handle the poll functions return 0, ffi_workflow_wait_status
returns 7 and ffi_workflow_destroy does nothing.
test "raw kernels" {
let input : FixedArray[Int64] = [5L, -2L, 7L, 0L]
assert_eq(@native.ffi_execute_reduce_min_i64(input, 4, 2, 2), -2L)
assert_eq(@native.ffi_execute_reduce_max_i64(input, 4, 2, 2), 7L)
let output : FixedArray[Int64] = FixedArray::make(4, 0L)
assert_eq(@native.ffi_execute_scan_sum_i64(input, 4, 2, 2, output), 0)
debug_inspect(output, content="<FixedArray: [5, 3, 10, 10]>")
}
Equality
T::equal
Structural equality, promoted on NativeMapRequest, NativeReduceRequest,
NativeScanRequest, NativeReductionKernel, NativeRequestError,
NativeValueType, WorkflowResult and WorkflowRuntimeState. Use == and
!=. WorkflowHandle has no equality.
pub fn WorkflowResult::equal(Self, Self) -> Bool