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 w≥⌈n/c⌉w \ge \lceil n / c \rceil 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 nn be the input length, ww the worker count and cc the chunk size. Every kernel requires

n>0,0<w≤n,0<c≤n,w≥⌈nc⌉,n > 0, \qquad 0 < w \le n, \qquad 0 < c \le n, \qquad w \ge \left\lceil \frac{n}{c} \right\rceil ,

and otherwise fails with InvalidArgument. With the defaults w=c=1w = c = 1 only an input of length one is accepted, so pass both arguments. The input is split into ⌈n/c⌉\lceil n / c \rceil 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 2xi2 x_i 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 x↦2xx \mapsto 2x 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 w=1,c=4w = 1, c = 4 but returns 0 with w=2,c=2w = 2, c = 2.

execute_scan_sum_i32

Returns the inclusive prefix sums yi=x0+⋯+xiy_i = x_0 + \dots + x_i.

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.