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

既定の実行ポリシーを返します。ネイティブバックエンド、同期モード、ワーカー 1、チャンクサイズ 1、入力順序の保持です。@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 で検証を通る 3 つのカーネルです。

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 の規則を満たしますが、ネイティブカーネルにはもう 1 つ条件があります。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 では長さ 1 の入力しか受け付けないので、両方の引数を渡してください。入力は ⌈n/c⌉\lceil n / c \rceil 個の連続したチャンクに分けられ、各チャンクの大きさの差は高々 1 です。これと以下の結果の導出はネイティブバックエンドの設計にあります。

execute_map_i32

すべての要素を 2 倍にして新しい配列を返します。ある 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 を返し、本当の和が 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

最初の 2 つの引数はノードの 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 が見つけた問題と、問題が 1 つもないかどうかを返します。

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 の項で述べた不具合のため、最後の例はドキュメントのテストとしてコンパイルされていません。