workflow チュートリアル

このチュートリアルでは、並行する作業をタスクグラフとして記述する方法を示します。ケイパビリティを宣言し、それを使うノードを追加し、エッジでノードを順序付け、何かが動く前に validate に結果を検査させます。このパッケージはすべてのターゲットで使えます。

クイックスタート

パッケージをインポートし、ポリシーや計算ノードが必要なら shared と plan もインポートします。

import {
  "Luna-Flow/luna_thread/plan",
  "Luna-Flow/luna_thread/shared",
  "Luna-Flow/luna_thread/workflow",
}

2 つのノードからなるフォークジョインのグラフです。

test "workflow quick start" {
  let graph = @workflow.Workflow::new("fork-join")
    .add_node(@workflow.spawn_node(1, "spawn"))
    .add_node(@workflow.join_node(2, "join"))
    .add_edge(@workflow.Edge::new(1, 2, @workflow.control_dependency()))
  inspect(@workflow.is_ready(graph), content="true")
}

日常的なタスク

チャネルで作業を受け渡す

通信するノードはケイパビリティを id で指定します。チャネルは MoveOnly アクセスで使い、send から recv へのエッジによって受信は送信を待ちます。

test "channel" {
  let jobs = @workflow.Capability::new(
    1,
    "jobs",
    @workflow.channel_capability(),
    @workflow.move_only_access(),
  )
  let graph = @workflow.Workflow::new("producer-consumer")
    .add_capability(jobs)
    .add_node(@workflow.send_node(1, "produce", capability=1))
    .add_node(@workflow.recv_node(2, "consume", capability=1))
    .add_edge(@workflow.Edge::new(1, 2, @workflow.data_dependency()))
  assert_eq(@workflow.validate(graph).length(), 0)
}

クリティカルセクションを守る

ミューテックスは SynchronizeOnly アクセスで、Lock ノードと Unlock ノードが使います。

test "critical section" {
  let guard_cap = @workflow.Capability::new(
    1,
    "guard",
    @workflow.mutex_capability(),
    @workflow.synchronize_only_access(),
  )
  let graph = @workflow.Workflow::new("critical")
    .add_capability(guard_cap)
    .add_node(@workflow.lock_node(1, "enter", capability=1))
    .add_node(@workflow.write_shared_node(2, "update", capability=2))
    .add_node(@workflow.unlock_node(3, "leave", capability=1))
    .add_edge(@workflow.Edge::new(1, 2, @workflow.synchronization_dependency()))
    .add_edge(@workflow.Edge::new(2, 3, @workflow.synchronization_dependency()))
  debug_inspect(@workflow.validate(graph), content="[MissingCapability(capability_id=2)]")
}

書き込みはケイパビリティ 2 を指定していますが、グラフはそれを宣言していません。ReadWrite アクセスの AtomicCell をケイパビリティ 2 として追加すれば直ります。

プランを埋め込む

計算ノードはプランを持ち、プラン自身の問題はノードの id と一緒に報告されます。

test "compute node" {
  let plan = @plan.map("halve", @plan.f64_type(), 8)
  let graph = @workflow.Workflow::new("compute").add_node(
    @workflow.compute_node(1, "halve", plan),
  )
  debug_inspect(
    @workflow.validate(graph),
    content="[InvalidComputePlan(node_id=1, issue=UnsupportedValueType)]",
  )
  assert_true(@workflow.validate(graph).any(@workflow.is_invalid_compute_plan))
}

ワークフローを投入する

submit はグラフを検証し、選んだバックエンドがそれを受理するかを記録します。

test "submit" {
  let graph = @workflow.Workflow::new("one").add_node(
    @workflow.spawn_node(1, "spawn"),
  )
  let submission = @workflow.submit(graph)
  inspect(submission.accepted(), content="true")
  inspect(submission.completed(), content="false")
}

さらに進んで

グラフを実行するには、ネイティブターゲットで @native.submit_workflow_async(またはファサードの submit_workflow_async)に渡します。やり方はネイティブバックエンドのチュートリアルに、どのノード種類がブロックし何がそれを起こすかはネイティブバックエンドの設計にあります。validate はすべての問題を返すので、ツールはグラフの問題を一度に示せます。述語 is_missing_capability、is_cyclic_dependency、is_invalid_compute_plan、is_unsupported_shared_write はよくある問題をまとめて判定します。

よくある落とし穴

  • add_capability、add_node、add_edge はワークフローをその場で変更します。同じワークフローに束縛された 2 つの変数からは同じノードが見えます。
  • 入ってくるエッジのないノードがあると、validate は 3 つ以上のノードからなる閉路を見逃します。ネイティブランタイムは投入時にそれを拒否します。
  • 相手より先に実行されうる Recv、Lock、Send は、ネイティブランタイムで永久にブロックすることがあります。そのようなノードはエッジで順序付けてください。
  • グループに Barrier ノードが 1 つしかないと、実行時にステータス 14 で失敗します。バリアには同じ深さに 2 つ以上のノードが必要です。
  • RwLock と Semaphore は validate を通りますが、ネイティブランタイムに拒否されます。

次のステップ

workflow API に validate のすべての検査が載っています。workflow の設計ではグラフのモデルと閉路の検査を説明しています。core チュートリアルではファサードからワークフローを作ります。