shared API

The package Luna-Flow/luna_thread/shared, imported as @shared, holds the vocabulary that every other package of the module uses: backend targets, execution modes, ordering guarantees, execution policies and their validation, the declared capabilities of each backend, native status codes, and integer-coded records that mirror the C request structs. It builds on every target.

The enumerations of this package are read-only outside it: match on their constructors, but build values with the functions below.

Targets, modes and orderings

BackendTarget

The runtime that executes a plan or workflow.

pub enum BackendTarget {
  Native
  JavaScript
} derive(Eq, @debug.Debug)

ExecutionMode

Whether a submission blocks until it finishes.

pub enum ExecutionMode {
  Synchronous
  Asynchronous
} derive(Eq, @debug.Debug)

OrderingGuarantee

Whether results must keep the order of the input.

pub enum OrderingGuarantee {
  PreserveInputOrder
  RelaxedOrder
} derive(Eq, @debug.Debug)

native_target, javascript_target, synchronous_mode, asynchronous_mode, preserve_input_order and relaxed_order

Return Native, JavaScript, Synchronous, Asynchronous, PreserveInputOrder and RelaxedOrder.

pub fn native_target() -> BackendTarget
pub fn javascript_target() -> BackendTarget
pub fn synchronous_mode() -> ExecutionMode
pub fn asynchronous_mode() -> ExecutionMode
pub fn preserve_input_order() -> OrderingGuarantee
pub fn relaxed_order() -> OrderingGuarantee

backend_label

Returns "native" or "javascript".

pub fn backend_label(BackendTarget) -> String
test "targets" {
  inspect(@shared.backend_label(@shared.javascript_target()), content="javascript")
  assert_true(@shared.relaxed_order() is @shared.RelaxedOrder)
}

Execution policies

ExecutionPolicy

How a plan or workflow should run: backend, mode, number of workers, chunk size and ordering.

pub struct ExecutionPolicy {
  backend : BackendTarget
  mode : ExecutionMode
  worker_count : Int
  chunk_size : Int
  ordering : OrderingGuarantee
} derive(Eq, @debug.Debug)

PolicyIssue and PolicyError

A reason why a policy is outside the v1 subset, and the error that carries the first such reason.

pub enum PolicyIssue {
  WorkerCountMustBePositive(Int)
  ChunkSizeMustBePositive(Int)
  UnsupportedBackendForV1(BackendTarget)
  UnsupportedModeForV1(ExecutionMode)
} derive(Eq, @debug.Debug)

pub struct PolicyError {
  issue : PolicyIssue
} derive(Eq, @debug.Debug)

make_execution_policy

Builds a policy and returns the first problem as a PolicyError.

pub fn make_execution_policy(backend? : BackendTarget, mode? : ExecutionMode, worker_count? : Int, chunk_size? : Int, ordering? : OrderingGuarantee) -> Result[ExecutionPolicy, PolicyError]

The defaults are Native, Synchronous, 1, 1 and PreserveInputOrder. The checks run in the order worker_count > 0, chunk_size > 0, backend == Native, mode == Synchronous; the ordering is not checked.

ExecutionPolicy::new

Builds a policy like make_execution_policy and aborts when it would return an error.

pub fn ExecutionPolicy::new(backend? : BackendTarget, mode? : ExecutionMode, worker_count? : Int, chunk_size? : Int, ordering? : OrderingGuarantee) -> Self

native_policy and javascript_policy

Return the default native policy, and attempt the same for the JavaScript backend.

pub fn native_policy() -> ExecutionPolicy
pub fn javascript_policy() -> ExecutionPolicy

native_policy() is make_execution_policy().unwrap(). javascript_policy() is make_execution_policy(backend=JavaScript).unwrap(), which aborts in v1 because the JavaScript backend is rejected.

validate_policy

Returns every issue of an existing policy, in the order of the checks above.

pub fn validate_policy(ExecutionPolicy) -> Array[PolicyIssue]

Policies built with the functions above always pass; validate_policy matters for a policy stored in a workflow, which @workflow.validate re-checks.

is_parallel

Returns true when the policy asks for more than one worker.

pub fn is_parallel(ExecutionPolicy) -> Bool
test "policies" {
  let policy = @shared.make_execution_policy(worker_count=4, chunk_size=8).unwrap()
  assert_true(@shared.is_parallel(policy))
  assert_eq(@shared.validate_policy(policy).length(), 0)
  let rejected = @shared.make_execution_policy(backend=@shared.javascript_target())
  debug_inspect(
    rejected,
    content="Err({ issue: UnsupportedBackendForV1(JavaScript) })",
  )
  assert_false(@shared.is_parallel(@shared.native_policy()))
}

Runtime capabilities

RuntimeCapabilities

The features a backend declares.

pub struct RuntimeCapabilities {
  backend : BackendTarget
  supports_parallelism : Bool
  supports_async : Bool
  supports_zero_copy_buffers : Bool
} derive(Eq, @debug.Debug)

RuntimeCapabilities::for_backend

Returns the declared table of a backend.

pub fn RuntimeCapabilities::for_backend(BackendTarget) -> Self
BackendParallelismAsyncZero-copy buffers
Nativetruefalsefalse
JavaScripttruetruetrue

The table is a fixed declaration, not a probe of the running system. The JavaScript entries describe the backend the specification plans; the JavaScript backend executes nothing yet.

Native status codes

NativeStatus

The status codes 0 to 7 of the C runtime.

pub enum NativeStatus {
  Ok
  InvalidArgument
  UnsupportedBackend
  UnsupportedMode
  UnsupportedValueType
  UnsupportedReductionKernel
  Overflow
  NullPointer
} derive(Eq, @debug.Debug)

native_status_code and native_status_from_code

Convert between NativeStatus and its C code.

pub fn native_status_code(NativeStatus) -> Int
pub fn native_status_from_code(Int) -> NativeStatus

native_status_code numbers the constructors 0 to 7 in declaration order. native_status_from_code inverts it on 0 to 7 and maps every other code, including the workflow statuses 8 to 14, to InvalidArgument.

test "status codes" {
  assert_true(@shared.native_status_from_code(6) is @shared.Overflow)
  assert_true(@shared.native_status_from_code(14) is @shared.InvalidArgument)
  for code in 0..<8 {
    let status = @shared.native_status_from_code(code)
    assert_eq(@shared.native_status_code(status), code)
  }
}

Integer-coded native records

These records mirror the C structs luna_thread_buffer, luna_thread_map_request, luna_thread_reduce_request and luna_thread_scan_request field by field, with enumerations as Int codes. No package of the module consumes them; backend/native uses its own typed records. They store their arguments unchecked.

NativeBuffer

A pointer, as an Int, and a length.

pub struct NativeBuffer {
  ptr : Int
  length : Int
} derive(Eq, @debug.Debug)
pub fn NativeBuffer::new(Int, Int) -> Self
pub fn NativeBuffer::ptr(Self) -> Int
pub fn NativeBuffer::length(Self) -> Int

NativeMapRequest

Input and output buffers, element count, value type code, worker count and chunk size.

pub struct NativeMapRequest {
  input : NativeBuffer
  output : NativeBuffer
  element_count : Int
  value_type : Int
  worker_count : Int
  chunk_size : Int
} derive(Eq, @debug.Debug)
pub fn NativeMapRequest::new(NativeBuffer, NativeBuffer, Int, Int, Int, Int) -> Self
pub fn NativeMapRequest::input(Self) -> NativeBuffer
pub fn NativeMapRequest::output(Self) -> NativeBuffer
pub fn NativeMapRequest::element_count(Self) -> Int
pub fn NativeMapRequest::value_type(Self) -> Int
pub fn NativeMapRequest::worker_count(Self) -> Int
pub fn NativeMapRequest::chunk_size(Self) -> Int

The constructor takes the fields in declaration order.

NativeReduceRequest and NativeScanRequest

The same fields as NativeMapRequest plus a reduction kernel code, placed after the value type.

pub struct NativeReduceRequest {
  input : NativeBuffer
  output : NativeBuffer
  element_count : Int
  value_type : Int
  reduction_kernel : Int
  worker_count : Int
  chunk_size : Int
} derive(Eq, @debug.Debug)
pub fn NativeReduceRequest::new(NativeBuffer, NativeBuffer, Int, Int, Int, Int, Int) -> Self
pub fn NativeReduceRequest::input(Self) -> NativeBuffer
pub fn NativeReduceRequest::output(Self) -> NativeBuffer
pub fn NativeReduceRequest::element_count(Self) -> Int
pub fn NativeReduceRequest::value_type(Self) -> Int
pub fn NativeReduceRequest::reduction_kernel(Self) -> Int
pub fn NativeReduceRequest::worker_count(Self) -> Int
pub fn NativeReduceRequest::chunk_size(Self) -> Int

pub struct NativeScanRequest {
  input : NativeBuffer
  output : NativeBuffer
  element_count : Int
  value_type : Int
  reduction_kernel : Int
  worker_count : Int
  chunk_size : Int
} derive(Eq, @debug.Debug)
pub fn NativeScanRequest::new(NativeBuffer, NativeBuffer, Int, Int, Int, Int, Int) -> Self
pub fn NativeScanRequest::input(Self) -> NativeBuffer
pub fn NativeScanRequest::output(Self) -> NativeBuffer
pub fn NativeScanRequest::element_count(Self) -> Int
pub fn NativeScanRequest::value_type(Self) -> Int
pub fn NativeScanRequest::reduction_kernel(Self) -> Int
pub fn NativeScanRequest::worker_count(Self) -> Int
pub fn NativeScanRequest::chunk_size(Self) -> Int
test "native records" {
  let buffer = @shared.NativeBuffer::new(0, 4)
  let request = @shared.NativeReduceRequest::new(buffer, buffer, 4, 0, 1, 2, 2)
  assert_eq(request.reduction_kernel(), 1)
  assert_eq(request.input().length(), 4)
}

Equality and package information

T::equal

Structural equality, promoted on every type of the package: BackendTarget, ExecutionMode, ExecutionPolicy, NativeBuffer, NativeMapRequest, NativeReduceRequest, NativeScanRequest, NativeStatus, OrderingGuarantee, PolicyError, PolicyIssue and RuntimeCapabilities. Use == and !=.

pub fn ExecutionPolicy::equal(Self, Self) -> Bool

package_name

Returns "shared".

pub fn package_name() -> String