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 下所述的覆盖条件 w≥⌈n/c⌉w \ge \lceil n / c \rceil。

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 运行时的一个内核,并阻塞直到完成。它们不使用计划。设 nn 为输入长度,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 失败。使用默认值 w=c=1w = c = 1 时只接受长度为一的输入,所以请两个参数都传。输入被分成 ⌈n/c⌉\lceil n / c \rceil 个连续的块,各块大小至多相差一;原生后端设计推导了这一点以及下面的结果。

execute_map_i32

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

pub fn execute_map_i32(FixedArray[Int], worker_count? : Int, chunk_size? : Int) -> Result[FixedArray[Int], @native.NativeRequestError]

输入不会被修改。在 v1 中映射固定为 x↦2xx \mapsto 2x。

execute_reduce_sum_i32

返回各元素之和。

pub fn execute_reduce_sum_i32(FixedArray[Int], worker_count? : Int, chunk_size? : Int) -> Int

返回的和总是精确的。任何失败(无效参数或部分和溢出)时,函数返回 0,这与真实的零和无法区分。部分和是否溢出取决于分块方式:[-1, 0, 2147483647, 1] 在 w=1,c=4w = 1, c = 4 时求得 2147483647,在 w=2,c=2w = 2, c = 2 时却返回 0。

execute_scan_sum_i32

返回包含式前缀和 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]

参数违反上述条件时以 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 下所述的缺陷,最后这个示例没有编译进文档测试。