core API
根包 Luna-Flow/luna_thread 以 @luna_thread 导入,是本模块的门面。它重新导出 plan、shared 和 workflow 中常用的构造函数并填好默认值,在原生后端上执行整数内核,并提交工作流。它只能在 native 目标上构建。这里的每个函数都是一层薄包装;所链接的包页面给出了其返回类型的完整语义。
模块信息
scaffold_status
返回固定的描述 "luna_thread MoonBit facade"。
pub fn scaffold_status() -> String
root_plan_package 和 root_workflow_package
返回 plan 和 workflow 包报告的自身名称,即 "plan" 和 "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")
}
策略与能力
default_backend
返回 Native,即未指定后端时计划、工作流和提交所使用的后端。
pub fn default_backend() -> @shared.BackendTarget
default_policy
返回默认执行策略:原生后端、同步模式、一个工作线程、块大小为一、保持输入顺序。它与 @shared.native_policy() 是同一个值。
pub fn default_policy() -> @shared.ExecutionPolicy
javascript_policy
为 JavaScript 后端构建策略。在 v1 中,策略构造器拒绝 Native 以外的所有后端,因此这个函数会中止;保留它是为了规范中描述的 JavaScript 后端。
pub fn javascript_policy() -> @shared.ExecutionPolicy
make_policy
构建一个执行策略并检查它,把第一个问题作为 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]
默认值依次为 Native、Synchronous、1、1 和 PreserveInputOrder。检查按以下顺序进行并返回第一个失败:worker_count > 0、chunk_size > 0、backend == Native、mode == Synchronous。结果与 @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
返回后端声明的能力表;参见 shared API 中的 RuntimeCapabilities::for_backend。
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)
}
值类型与归约内核
i32_type、i64_type、f32_type、f64_type、bytes_type 和 opaque_type
返回计划输入域的元素类型 I32、I64、F32、F64、Bytes 和 Opaque(name)。在 v1 中只有 I32 和 I64 能通过验证。
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 和 max_reduction
返回归约内核 Sum、Min 和 Max,即在 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")
}
计划
map
构建一个作用于 input_length 个同类型元素的映射计划。
pub fn map(String, @plan.ValueType, Int, policy? : @shared.ExecutionPolicy, ordering? : @shared.OrderingGuarantee) -> @plan.Plan
参数依次是标签、元素类型和输入长度。策略默认为 default_policy(),顺序默认为 PreserveInputOrder。计划不会被检查;请调用 validate 或 is_ready。
reduce
构建一个用归约内核合并输入的归约计划。
pub fn reduce(String, @plan.ValueType, Int, @plan.ReductionKernel, policy? : @shared.ExecutionPolicy) -> @plan.Plan
归约计划的顺序总是 PreserveInputOrder。
scan
构建一个扫描(前缀)计划。
pub fn scan(String, @plan.ValueType, Int, policy? : @shared.ExecutionPolicy) -> @plan.Plan
扫描计划中没有归约内核,并且总是保持输入顺序。
map_reduce
构建一个先映射每个元素、再归约结果的计划。
pub fn map_reduce(String, @plan.ValueType, Int, @plan.ReductionKernel, policy? : @shared.ExecutionPolicy, ordering? : @shared.OrderingGuarantee) -> @plan.Plan
validate
返回计划超出 v1 子集的所有原因,若没有则返回空数组。
pub fn validate(@plan.Plan) -> Array[@plan.ValidationIssue]
这就是 @plan.validate;plan API 列出了各项检查。
is_ready
当 validate 没有发现问题时返回 true。
pub fn is_ready(@plan.Plan) -> Bool
就绪的计划满足 v1 的规则,但原生内核还有一个要求,即 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]")
}
直接执行
这些函数在 FixedArray[Int] 上运行原生 C 运行时的一个内核,并阻塞直到完成。它们不使用计划。设 为输入长度, 为工作线程数, 为块大小。每个内核都要求
否则以 InvalidArgument 失败。使用默认值 时只接受长度为一的输入,所以请两个参数都传。输入被分成 个连续的块,各块大小至多相差一;原生后端设计推导了这一点以及下面的结果。
execute_map_i32
把每个元素翻倍并返回新数组;若某个 超出 32 位范围,则返回 Overflow。
pub fn execute_map_i32(FixedArray[Int], worker_count? : Int, chunk_size? : Int) -> Result[FixedArray[Int], @native.NativeRequestError]
输入不会被修改。在 v1 中映射固定为 。
execute_reduce_sum_i32
返回各元素之和。
pub fn execute_reduce_sum_i32(FixedArray[Int], worker_count? : Int, chunk_size? : Int) -> Int
返回的和总是精确的。任何失败(无效参数或部分和溢出)时,函数返回 0,这与真实的零和无法区分。部分和是否溢出取决于分块方式:[-1, 0, 2147483647, 1] 在 时求得 2147483647,在 时却返回 0。
execute_scan_sum_i32
返回包含式前缀和 。
pub fn execute_scan_sum_i32(FixedArray[Int], worker_count? : Int, chunk_size? : Int) -> Result[FixedArray[Int], @native.NativeRequestError]
参数违反上述条件时以 InvalidArgument 失败,部分和溢出时以 Overflow 失败。
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)")
}
工作流
workflow
创建一个带标签和策略(默认为 default_policy())的空工作流。
pub fn workflow(String, policy? : @shared.ExecutionPolicy) -> @workflow.Workflow
使用 workflow API 中的 Workflow::add_capability、Workflow::add_node 和 Workflow::add_edge 添加能力、节点和边。这些方法会就地修改工作流。
compute_task、spawn_task 和 join_task
分别创建一个携带计划的计算节点、一个派生节点和一个汇合节点。
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
前两个参数是节点 id 和标签。这些节点都不需要能力。
channel_capability、mutex_capability 和 shared_read_capability
以各自种类唯一有效的访问模式创建能力:通道使用 MoveOnly,互斥锁使用 SynchronizeOnly,共享只读视图使用 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 和 workflow_is_ready
返回 @workflow.validate 找到的问题,以及是否一个问题都没有。
pub fn workflow_validate(@workflow.Workflow) -> Array[@workflow.WorkflowIssue]
pub fn workflow_is_ready(@workflow.Workflow) -> Bool
submit_workflow
针对某个后端验证工作流,并把结果记录为 Submission。
pub fn submit_workflow(@workflow.Workflow, backend? : @shared.BackendTarget) -> @workflow.Submission
当工作流没有问题且后端为 Native 时,提交被接受。不会执行任何东西:completed 总是 false。要运行工作流,请使用 submit_workflow_async。
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())
}
异步工作流
submit_workflow_async
在原生运行时上启动一个工作流,并立即返回句柄。
pub fn submit_workflow_async(@workflow.Workflow) -> Result[@native.WorkflowHandle, @native.NativeRequestError]
工作流不会先在 MoonBit 中验证,而是由 C 运行时检查。当前实现总是返回 Ok。当 C 运行时拒绝该图时,句柄为空,wait_workflow 报告状态 Submitted、状态码 7(NullPointer)。
poll_workflow
不等待,返回正在运行的工作流的快照。
pub fn poll_workflow(@native.WorkflowHandle) -> @native.WorkflowResult
wait_workflow
阻塞直到工作流完成或失败,然后返回其结果。
pub fn wait_workflow(@native.WorkflowHandle) -> @native.WorkflowResult
成功时 status 为 0,否则为运行时状态码,例如 14(BARRIER_BROKEN);completed_nodes 统计已完成的节点,failed_node_id 指出失败的节点,或为 -1。若节点阻塞后从未被唤醒,工作流永远不会结束,wait_workflow 也不会返回;参见原生后端设计。
drop_workflow
停止工作流的工作线程并释放其运行时。
pub fn drop_workflow(@native.WorkflowHandle) -> Unit
每个句柄恰好调用一次,且在最后一次 poll_workflow 或 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 }
由于 submit_workflow_async 下所述的缺陷,最后这个示例没有编译进文档测试。