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,因此在该构建中内核运行在调用线程上;参见原生后端设计。
内核
所有内核都接受一个长度为 的 FixedArray[Int]、工作线程数 和块大小 ,阻塞直到完成,并且不修改输入。它们要求
否则以 InvalidArgument 失败。core API 中同名的门面函数用可选参数包装了这些函数。
execute_map_i32
返回一个每个元素都翻倍的新数组;若某个 超出 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