workflow 教程

本教程展示如何把并发工作描述为任务图:声明能力,添加使用它们的节点,用边为节点排序,并在任何东西运行之前让 validate 检查结果。这个包可以在所有目标上使用。

快速开始

导入该包;需要策略和计算节点时,再导入 shared 和 plan:

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

一个只有两个节点的派生-汇合图:

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 会就地修改工作流。绑定到同一个工作流的两个变量看到的是同样的节点。
  • 当某个节点没有入边时,validate 会漏掉三个或更多节点构成的环;原生运行时会在提交时拒绝它们。
  • 可能先于其配对节点运行的 Recv、Lock 或 Send 可能会在原生运行时上永远阻塞。请用边为这类节点排序。
  • 组内只有一个 Barrier 节点时,会在运行时以状态码 14 失败;一个屏障需要在同一深度上至少有两个节点。
  • RwLock 和 Semaphore 能通过 validate,但会被原生运行时拒绝。

下一步

workflow API 列出了 validate 的每项检查。workflow 设计给出了图模型和环检查。core 教程通过门面构建工作流。