backend/native API

包 Luna-Flow/luna_thread/backend/native 以 @native 导入,是基于 C 外部函数接口的后端。它在 MoonBit 数组上运行 C 运行时的整数内核,把工作流提交给 C 调度器,并把计划转换为带类型的原生请求。它只能在 native 目标上构建,并通过原生桩链接 C 运行时。

本包的枚举类型在包外是只读的:可以对其构造器做模式匹配。请求记录由 *_from_plan 函数构建。

后端信息

backend_target

返回 Native。

pub fn backend_target() -> @shared.BackendTarget

default_plan_kind

返回 Map。

pub fn default_plan_kind() -> @plan.PlanKind

package_name

返回 "backend/native"。

pub fn package_name() -> String

supports_openmp

返回 true。

pub fn supports_openmp() -> Bool

这个值是常量。moon 构建在编译 C 桩时没有启用 OpenMP,因此在该构建中内核运行在调用线程上;参见原生后端设计。

内核

所有内核都接受一个长度为 nn 的 FixedArray[Int]、工作线程数 ww 和块大小 cc,阻塞直到完成,并且不修改输入。它们要求

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

否则以 InvalidArgument 失败。core API 中同名的门面函数用可选参数包装了这些函数。

execute_map_i32

返回一个每个元素都翻倍的新数组;若某个 2xi2 x_i 超出 32 位范围,则返回 Overflow。

pub fn execute_map_i32(FixedArray[Int], Int, Int) -> Result[FixedArray[Int], NativeRequestError]

execute_reduce_sum_i32

返回各元素之和,任何失败时返回 0。

pub fn execute_reduce_sum_i32(FixedArray[Int], Int, Int) -> Int

内核先从左到右累加每个块,再从左到右累加各块的和;只要其中任何一个部分和溢出就失败。返回的非零值就是精确的和;0 则既可能是和,也可能表示失败。

execute_scan_sum_i32

返回包含式前缀和。

pub fn execute_scan_sum_i32(FixedArray[Int], Int, Int) -> Result[FixedArray[Int], NativeRequestError]

若某个部分和溢出,则以 Overflow 失败。

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)")
}

工作流

workflow_is_supported

当 @workflow.validate 没有发现任何问题,且工作流的策略指定原生后端时,返回 true。

pub fn workflow_is_supported(@workflow.Workflow) -> Bool

submit_workflow

返回 @workflow.submit(workflow, backend=Native):一条经过验证的提交记录。它不会运行工作流。

pub fn submit_workflow(@workflow.Workflow) -> @workflow.Submission

WorkflowRuntimeState

工作流在 C 运行时中的状态。

pub enum WorkflowRuntimeState {
  Submitted
  Running
  Completed
  Failed
  Rejected
} derive(Eq, @debug.Debug)

在某个工作线程取走第一个节点之前,状态一直是 Submitted。运行时从不设置 Rejected;被拒绝的提交会得到一个空句柄。

WorkflowResult

工作流的一个快照:状态、运行时状态码、已完成节点数,以及失败节点的 id(没有失败时为 -1)。

pub struct WorkflowResult {
  state : WorkflowRuntimeState
  status : Int
  completed_nodes : Int
  failed_node_id : Int
} derive(Eq, @debug.Debug)

status 是架构指南中列出的状态码之一。

NativeWorkflowHandle 和 WorkflowHandle

指向 C 运行时中正在运行的工作流的不透明指针,及其 MoonBit 包装。

#external
pub type NativeWorkflowHandle

pub struct WorkflowHandle {
  raw : NativeWorkflowHandle
}

submit_workflow_async

把工作流复制为扁平的整数数组,并以其策略中的 worker_count 个线程在 C 调度器上启动它。

pub fn submit_workflow_async(@workflow.Workflow) -> Result[WorkflowHandle, NativeRequestError]

工作流不会在 MoonBit 中验证;由 C 运行时检查,若拒绝该图则返回空句柄。该函数总是返回 Ok。计算节点会被调度,但其计划不会被执行。

poll_workflow

不等待,直接返回一个快照。

pub fn poll_workflow(WorkflowHandle) -> WorkflowResult

wait_workflow

等待工作流进入 Completed 或 Failed,然后返回快照。

pub fn wait_workflow(WorkflowHandle) -> WorkflowResult

对空句柄,它会立即返回,状态为 Submitted,状态码为 7。如果工作流发生死锁,它不会返回。

drop_workflow

停止并汇合运行时线程,然后释放工作流。每个句柄只调用一次;对空句柄不做任何事。

pub fn drop_workflow(WorkflowHandle) -> Unit

带类型的原生请求

NativeValueType 和 NativeReductionKernel

C 运行时实现的值类型和归约内核。

pub enum NativeValueType {
  I32
  I64
} derive(Eq, @debug.Debug)

pub enum NativeReductionKernel {
  Sum
  Min
  Max
} derive(Eq, @debug.Debug)

NativeRequestError

计划无法转换为原生请求的原因,或内核失败的原因。

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

NativeMapRequest、NativeReduceRequest 和 NativeScanRequest

C 请求结构体的带类型对应物。

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)

由计划构建的请求,其 input 和 output 缓冲区都是 NativeBuffer::new(0, 0):计划描述的是数据的形状,而不是数据本身。本模块中没有任何函数把这些记录发送给 C;它们是经过验证的描述,将来的执行器会为其填入缓冲区。

map_request_from_plan

把 Map 计划转换为 NativeMapRequest。

pub fn map_request_from_plan(@plan.Plan) -> Result[NativeMapRequest, NativeRequestError]

若计划不是 Map 计划,则以 UnsupportedMode 失败;否则以 @plan.validate 的第一个问题所映射的错误失败:值类型、归约内核、后端和模式的问题映射为同名错误,MissingReductionKernel 映射为 UnsupportedReductionKernel,其余问题映射为 InvalidArgument。

reduce_request_from_plan

把 Reduce 或 MapReduce 计划转换为 NativeReduceRequest,错误规则相同。

pub fn reduce_request_from_plan(@plan.Plan) -> Result[NativeReduceRequest, NativeRequestError]

scan_request_from_plan

把 Scan 计划转换为内核为 Sum 的 NativeScanRequest,错误规则相同。

pub fn scan_request_from_plan(@plan.Plan) -> Result[NativeScanRequest, NativeRequestError]

map_request_is_valid、reduce_request_is_valid 和 scan_request_is_valid

检查 C 验证器同样会检查的请求字段:元素数、工作线程数和块大小为正,值类型受支持,对扫描还要求内核为 Sum。

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 和 scan_request_kernel

返回请求的值类型(I32 为 0,I64 为 1)和内核(Sum 为 0,Min 为 1,Max 为 2)的 C 编码。

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)")
}

原始外部函数

这些 extern "C" 声明是公开的,使调用者可以使用带类型函数没有包装的内核,例如 Int64 内核以及最小值和最大值归约。数组在调用期间被借用。映射和扫描返回状态码,并写入 output,它必须至少有 length 个元素;归约直接返回结果,失败时返回 0。参数依次是输入、其长度、工作线程数和块大小。

ffi_execute_map_i32 和 ffi_execute_map_i64

把每个元素翻倍后写入 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 和 ffi_execute_reduce_max_i64

返回输入的和、最小值或最大值。

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

最小值和最大值不会溢出。

ffi_execute_scan_sum_i32 和 ffi_execute_scan_sum_i64

把包含式前缀和写入 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

从扁平数组启动一个工作流:工作线程数;能力的 id、种类编码和数量;节点的 id、种类编码、能力 id(无则为 -1)和数量;边的起点、终点、种类编码和数量。运行时拒绝该图时返回空句柄。

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 和 ffi_workflow_destroy

读取工作流快照的一个字段、等待并返回最终状态,或销毁工作流。

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

对空句柄,轮询函数返回 0,ffi_workflow_wait_status 返回 7,ffi_workflow_destroy 不做任何事。

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]>")
}

相等性

T::equal

结构相等,作为方法提升到 NativeMapRequest、NativeReduceRequest、NativeScanRequest、NativeReductionKernel、NativeRequestError、NativeValueType、WorkflowResult 和 WorkflowRuntimeState 上。请使用 == 和 !=。WorkflowHandle 没有相等性。

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