core API
The root package Luna-Flow/luna_thread, imported as @luna_thread, is the
facade of the module. It re-exports the common constructors of plan,
shared and workflow with defaults filled in, executes integer kernels on
the native backend, and submits workflows. It builds only for the native
target. Every function here is a thin wrapper; the linked package pages give
the full semantics of the types it returns.
Module information
scaffold_status
Returns the fixed description "luna_thread MoonBit facade".
pub fn scaffold_status() -> String
root_plan_package and root_workflow_package
Return the names that the plan and workflow packages report for
themselves, "plan" and "workflow".
pub fn root_plan_package() -> String
pub fn root_workflow_package() -> String
test "module information" {
inspect(@luna_thread.scaffold_status(), content="luna_thread MoonBit facade")
inspect(@luna_thread.root_plan_package(), content="plan")
inspect(@luna_thread.root_workflow_package(), content="workflow")
}
Policies and capabilities
default_backend
Returns Native, the backend that plans, workflows and submissions use when
none is given.
pub fn default_backend() -> @shared.BackendTarget
default_policy
Returns the default execution policy: native backend, synchronous mode, one
worker, chunk size one, input order preserved. It is the same value as
@shared.native_policy().
pub fn default_policy() -> @shared.ExecutionPolicy
javascript_policy
Builds a policy for the JavaScript backend. In v1 the policy constructor
rejects every backend other than Native, so this function aborts; it is kept
for the JavaScript backend that the specification describes.
pub fn javascript_policy() -> @shared.ExecutionPolicy
make_policy
Builds an execution policy and checks it, returning the first problem as a
PolicyError.
pub fn make_policy(backend? : @shared.BackendTarget, mode? : @shared.ExecutionMode, worker_count? : Int, chunk_size? : Int, ordering? : @shared.OrderingGuarantee) -> Result[@shared.ExecutionPolicy, @shared.PolicyError]
The defaults are Native, Synchronous, 1, 1 and PreserveInputOrder.
The checks run in this order and the first failure is returned:
worker_count > 0, chunk_size > 0, backend == Native,
mode == Synchronous. The result is the same as
@shared.make_execution_policy.
test "make_policy" {
let policy = @luna_thread.make_policy(worker_count=4, chunk_size=2).unwrap()
assert_eq(policy.worker_count, 4)
guard @luna_thread.make_policy(worker_count=0) is Err(error) else {
fail("expected an error")
}
debug_inspect(error.issue, content="WorkerCountMustBePositive(0)")
}
runtime_capabilities
Returns the declared capability table of a backend; see
RuntimeCapabilities::for_backend in the shared API.
pub fn runtime_capabilities(@shared.BackendTarget) -> @shared.RuntimeCapabilities
test "runtime_capabilities" {
let native = @luna_thread.runtime_capabilities(@luna_thread.default_backend())
assert_true(native.supports_parallelism)
assert_true(!native.supports_async)
}
Value types and reduction kernels
i32_type, i64_type, f32_type, f64_type, bytes_type and opaque_type
Return the element types I32, I64, F32, F64, Bytes and
Opaque(name) of a plan’s input domain. Only I32 and I64 pass validation
in v1.
pub fn i32_type() -> @plan.ValueType
pub fn i64_type() -> @plan.ValueType
pub fn f32_type() -> @plan.ValueType
pub fn f64_type() -> @plan.ValueType
pub fn bytes_type() -> @plan.ValueType
pub fn opaque_type(String) -> @plan.ValueType
sum_reduction, min_reduction and max_reduction
Return the reduction kernels Sum, Min and Max, the three kernels that
pass validation in v1.
pub fn sum_reduction() -> @plan.ReductionKernel
pub fn min_reduction() -> @plan.ReductionKernel
pub fn max_reduction() -> @plan.ReductionKernel
test "value types and kernels" {
debug_inspect(@luna_thread.i64_type(), content="I64")
debug_inspect(@luna_thread.opaque_type("rgba"), content="Opaque(\"rgba\")")
debug_inspect(@luna_thread.max_reduction(), content="Max")
}
Plans
map
Builds a map plan over input_length elements of one type.
pub fn map(String, @plan.ValueType, Int, policy? : @shared.ExecutionPolicy, ordering? : @shared.OrderingGuarantee) -> @plan.Plan
The arguments are the label, the element type and the input length. The
policy defaults to default_policy() and the ordering to
PreserveInputOrder. The plan is not checked; call validate or is_ready.
reduce
Builds a reduce plan that combines the input with a reduction kernel.
pub fn reduce(String, @plan.ValueType, Int, @plan.ReductionKernel, policy? : @shared.ExecutionPolicy) -> @plan.Plan
The ordering of a reduce plan is always PreserveInputOrder.
scan
Builds a scan (prefix) plan.
pub fn scan(String, @plan.ValueType, Int, policy? : @shared.ExecutionPolicy) -> @plan.Plan
A scan has no reduction kernel in the plan and always preserves input order.
map_reduce
Builds a plan that maps every element and then reduces the results.
pub fn map_reduce(String, @plan.ValueType, Int, @plan.ReductionKernel, policy? : @shared.ExecutionPolicy, ordering? : @shared.OrderingGuarantee) -> @plan.Plan
validate
Returns every reason why a plan is outside the v1 subset, or an empty array.
pub fn validate(@plan.Plan) -> Array[@plan.ValidationIssue]
This is @plan.validate; the plan API lists the checks.
is_ready
Returns true when validate finds no issue.
pub fn is_ready(@plan.Plan) -> Bool
A ready plan satisfies the v1 rules, but the native kernels have one more
requirement, the cover condition described under
execute_map_i32.
test "plans" {
let policy = @luna_thread.make_policy(worker_count=2, chunk_size=4).unwrap()
let sum = @luna_thread.map_reduce(
"sum-of-doubles",
@luna_thread.i32_type(),
8,
@luna_thread.sum_reduction(),
policy~,
)
assert_true(@luna_thread.is_ready(sum))
let floats = @luna_thread.map("halve", @luna_thread.f64_type(), 8, policy~)
debug_inspect(@luna_thread.validate(floats), content="[UnsupportedValueType]")
}
Direct execution
These functions run a kernel of the native C runtime on a FixedArray[Int]
and block until it finishes. They do not use plans. Let be the input
length, the worker count and the chunk size. Every kernel requires
and otherwise fails with InvalidArgument. With the defaults only
an input of length one is accepted, so pass both arguments. The input is split
into contiguous chunks whose sizes differ by at most one;
the native backend design derives this and the
results below.
execute_map_i32
Doubles every element, returning a new array, or Overflow when some
does not fit in 32 bits.
pub fn execute_map_i32(FixedArray[Int], worker_count? : Int, chunk_size? : Int) -> Result[FixedArray[Int], @native.NativeRequestError]
The input is not modified. The map is fixed to in v1.
execute_reduce_sum_i32
Returns the sum of the elements.
pub fn execute_reduce_sum_i32(FixedArray[Int], worker_count? : Int, chunk_size? : Int) -> Int
The sum is exact when it is returned. On any failure, an invalid argument or
an overflow of a partial sum, the function returns 0, which cannot be told
apart from a true sum of zero. Whether a partial sum overflows depends on the
chunking: [-1, 0, 2147483647, 1] sums to 2147483647 with
but returns 0 with .
execute_scan_sum_i32
Returns the inclusive prefix sums .
pub fn execute_scan_sum_i32(FixedArray[Int], worker_count? : Int, chunk_size? : Int) -> Result[FixedArray[Int], @native.NativeRequestError]
It fails with InvalidArgument when the arguments break the conditions above
and with Overflow when a partial sum overflows.
test "direct execution" {
let input : FixedArray[Int] = [1, 2, 3, 4, 5]
let doubled = @luna_thread.execute_map_i32(input, worker_count=2, chunk_size=3)
debug_inspect(doubled, content="Ok(<FixedArray: [2, 4, 6, 8, 10]>)")
let total = @luna_thread.execute_reduce_sum_i32(
input,
worker_count=2,
chunk_size=3,
)
assert_eq(total, 15)
let prefix = @luna_thread.execute_scan_sum_i32(input, worker_count=2, chunk_size=3)
debug_inspect(prefix, content="Ok(<FixedArray: [1, 3, 6, 10, 15]>)")
let rejected = @luna_thread.execute_map_i32(input)
debug_inspect(rejected, content="Err(InvalidArgument)")
}
Workflows
workflow
Creates an empty workflow with a label and a policy (default
default_policy()).
pub fn workflow(String, policy? : @shared.ExecutionPolicy) -> @workflow.Workflow
Add capabilities, nodes and edges with Workflow::add_capability,
Workflow::add_node and Workflow::add_edge from the
workflow API. These methods change the workflow in place.
compute_task, spawn_task and join_task
Create a compute node that carries a plan, a spawn node and a join node.
pub fn compute_task(Int, String, @plan.Plan) -> @workflow.Node
pub fn spawn_task(Int, String) -> @workflow.Node
pub fn join_task(Int, String) -> @workflow.Node
The first two arguments are the node id and label. None of these nodes needs a capability.
channel_capability, mutex_capability and shared_read_capability
Create capabilities with the only access mode that is valid for their kind:
a channel with MoveOnly, a mutex with SynchronizeOnly and a shared read
view with ReadOnly.
pub fn channel_capability(Int, String) -> @workflow.Capability
pub fn mutex_capability(Int, String) -> @workflow.Capability
pub fn shared_read_capability(Int, String) -> @workflow.Capability
workflow_validate and workflow_is_ready
Return the issues found by @workflow.validate, and whether there are none.
pub fn workflow_validate(@workflow.Workflow) -> Array[@workflow.WorkflowIssue]
pub fn workflow_is_ready(@workflow.Workflow) -> Bool
submit_workflow
Validates a workflow for a backend and records the outcome as a Submission.
pub fn submit_workflow(@workflow.Workflow, backend? : @shared.BackendTarget) -> @workflow.Submission
The submission is accepted when the workflow has no issues and the backend is
Native. Nothing is executed: completed is always false. Use
submit_workflow_async to run a workflow.
test "workflows" {
let policy = @luna_thread.make_policy(worker_count=2, chunk_size=2).unwrap()
let step = @luna_thread.map("double", @luna_thread.i32_type(), 4, policy~)
let graph = @luna_thread.workflow("pipeline", policy~)
.add_capability(@luna_thread.channel_capability(1, "jobs"))
.add_node(@luna_thread.spawn_task(1, "spawn"))
.add_node(@luna_thread.compute_task(2, "double", step))
.add_node(@luna_thread.join_task(3, "join"))
.add_edge(@workflow.Edge::new(1, 2, @workflow.control_dependency()))
.add_edge(@workflow.Edge::new(2, 3, @workflow.control_dependency()))
assert_true(@luna_thread.workflow_is_ready(graph))
let submission = @luna_thread.submit_workflow(graph)
assert_true(submission.accepted())
assert_true(!submission.completed())
}
Asynchronous workflows
submit_workflow_async
Starts a workflow on the native runtime and returns a handle at once.
pub fn submit_workflow_async(@workflow.Workflow) -> Result[@native.WorkflowHandle, @native.NativeRequestError]
The workflow is not validated in MoonBit first; the C runtime checks it. The
current implementation always returns Ok. When the C runtime rejects the
graph, the handle is empty and wait_workflow reports state Submitted with
status 7 (NullPointer).
poll_workflow
Returns a snapshot of a running workflow without waiting.
pub fn poll_workflow(@native.WorkflowHandle) -> @native.WorkflowResult
wait_workflow
Blocks until the workflow completes or fails, then returns its result.
pub fn wait_workflow(@native.WorkflowHandle) -> @native.WorkflowResult
status is 0 on success or a runtime status code such as 14
(BARRIER_BROKEN); completed_nodes counts finished nodes and
failed_node_id names the failing node, or is -1. A workflow whose nodes
block without being woken never finishes, and wait_workflow does not return;
see the native backend design.
drop_workflow
Stops the worker threads of a workflow and frees its runtime.
pub fn drop_workflow(@native.WorkflowHandle) -> Unit
Call it exactly once per handle, after the last poll_workflow or
wait_workflow.
let policy = @luna_thread.make_policy(worker_count=2, chunk_size=1).unwrap()
let graph = @luna_thread.workflow("fork-join", policy~)
.add_node(@luna_thread.spawn_task(1, "spawn"))
.add_node(@luna_thread.join_task(2, "join"))
.add_edge(@workflow.Edge::new(1, 2, @workflow.control_dependency()))
guard @luna_thread.submit_workflow_async(graph) is Ok(handle) else { return }
let result = @luna_thread.wait_workflow(handle)
@luna_thread.drop_workflow(handle)
// result: { state: Completed, status: 0, completed_nodes: 2, failed_node_id: -1 }
The last example is not compiled into the documentation tests because of the
defect described under submit_workflow_async.