workflow API

包 Luna-Flow/luna_thread/workflow 以 @workflow 导入,用于描述任务图:一组能力(通道、锁、屏障和共享数据)、一组进行计算或同步的节点,以及节点之间的依赖边。它验证图并记录提交。它本身不执行任何东西,可以在所有目标上构建;工作流由 backend/native 运行。

本包的枚举类型在包外是只读的:可以对其构造器做模式匹配,但请用下面的函数构建值。

能力

CapabilityKind

能力所代表的资源种类。

pub enum CapabilityKind {
  OwnedBuffer
  SharedReadView
  AtomicCell
  Mutex
  Condvar
  RwLock
  Semaphore
  Barrier
  Channel
  Opaque(String)
} derive(Eq, @debug.Debug)

Opaque(name) 永远不受支持。原生运行时还会拒绝 RwLock 和 Semaphore。

AccessMode

节点可以如何使用能力。

pub enum AccessMode {
  ReadOnly
  WriteOnly
  ReadWrite
  SynchronizeOnly
  MoveOnly
} derive(Eq, @debug.Debug)

每种能力只允许部分访问模式:

种类有效的访问模式
OwnedBufferReadOnly, WriteOnly, ReadWrite
SharedReadViewReadOnly
AtomicCellReadWrite
Mutex, Condvar, RwLock, Semaphore, BarrierSynchronizeOnly
ChannelMoveOnly
Opaque(_)无

Capability

带有 id、标签、种类和访问模式的能力。

pub struct Capability {
  id : Int
  label : String
  kind : CapabilityKind
  access : AccessMode
} derive(Eq, @debug.Debug)
pub fn Capability::new(Int, String, CapabilityKind, AccessMode) -> Self

Capability::new 不检查访问模式;由 validate 检查。

owned_buffer_capability、shared_read_view_capability、atomic_cell_capability、mutex_capability、condvar_capability、rwlock_capability、semaphore_capability、barrier_capability、channel_capability 和 opaque_capability

返回同名的能力种类。

pub fn owned_buffer_capability() -> CapabilityKind
pub fn shared_read_view_capability() -> CapabilityKind
pub fn atomic_cell_capability() -> CapabilityKind
pub fn mutex_capability() -> CapabilityKind
pub fn condvar_capability() -> CapabilityKind
pub fn rwlock_capability() -> CapabilityKind
pub fn semaphore_capability() -> CapabilityKind
pub fn barrier_capability() -> CapabilityKind
pub fn channel_capability() -> CapabilityKind
pub fn opaque_capability(String) -> CapabilityKind

read_only_access、write_only_access、read_write_access、synchronize_only_access 和 move_only_access

返回同名的访问模式。

pub fn read_only_access() -> AccessMode
pub fn write_only_access() -> AccessMode
pub fn read_write_access() -> AccessMode
pub fn synchronize_only_access() -> AccessMode
pub fn move_only_access() -> AccessMode
test "capabilities" {
  let lock = @workflow.Capability::new(
    1,
    "guard",
    @workflow.mutex_capability(),
    @workflow.synchronize_only_access(),
  )
  assert_true(lock.kind is @workflow.Mutex)
  assert_true(lock.access is @workflow.SynchronizeOnly)
}

节点

NodeKind

节点做什么。

pub enum NodeKind {
  Compute(@plan.Plan)
  Spawn
  Join
  Send
  Recv
  Lock
  Unlock
  Wait
  Signal
  Barrier
  ReadShared
  WriteShared
} derive(Eq, @debug.Debug)

Compute、Spawn 和 Join 不需要能力。其他种类必须指定一个种类匹配的能力:

节点种类能力种类
Send, RecvChannel
Lock, UnlockMutex, RwLock
Wait, SignalCondvar
BarrierBarrier
ReadSharedSharedReadView, AtomicCell
WriteSharedAtomicCell、OwnedBuffer(访问模式不能是 ReadOnly)

Node

带有 id、标签、种类和可选能力 id 的节点。

pub struct Node {
  id : Int
  label : String
  kind : NodeKind
  capability : Int?
} derive(Eq, @debug.Debug)
pub fn Node::new(Int, String, NodeKind, capability? : Int) -> Self

compute_node

创建一个不带能力的 Compute(plan) 节点。

pub fn compute_node(Int, String, @plan.Plan) -> Node

spawn_node、join_node、send_node、recv_node、lock_node、unlock_node、wait_node、signal_node、barrier_node、read_shared_node 和 write_shared_node

创建同名种类的节点,参数为 id、标签和可选的能力 id。

pub fn spawn_node(Int, String, capability? : Int) -> Node
pub fn join_node(Int, String, capability? : Int) -> Node
pub fn send_node(Int, String, capability? : Int) -> Node
pub fn recv_node(Int, String, capability? : Int) -> Node
pub fn lock_node(Int, String, capability? : Int) -> Node
pub fn unlock_node(Int, String, capability? : Int) -> Node
pub fn wait_node(Int, String, capability? : Int) -> Node
pub fn signal_node(Int, String, capability? : Int) -> Node
pub fn barrier_node(Int, String, capability? : Int) -> Node
pub fn read_shared_node(Int, String, capability? : Int) -> Node
pub fn write_shared_node(Int, String, capability? : Int) -> Node
test "nodes" {
  let send = @workflow.send_node(2, "send", capability=1)
  assert_eq(send.capability, Some(1))
  assert_true(send.kind is @workflow.Send)
  assert_eq(@workflow.join_node(3, "join").capability, None)
}

边

EdgeKind

为什么一个节点必须在另一个之前运行。

pub enum EdgeKind {
  DataDependency
  ControlDependency
  OwnershipTransfer
  SynchronizationDependency
} derive(Eq, @debug.Debug)

四种边以相同的方式对端点排序;种类只说明原因,并原样传给运行时。

Edge

从节点 from 到节点 to 的依赖:to 只能在 from 完成后开始。

pub struct Edge {
  from : Int
  to : Int
  kind : EdgeKind
} derive(Eq, @debug.Debug)
pub fn Edge::new(Int, Int, EdgeKind) -> Self

data_dependency、control_dependency、ownership_transfer_dependency 和 synchronization_dependency

返回 DataDependency、ControlDependency、OwnershipTransfer 和 SynchronizationDependency。

pub fn data_dependency() -> EdgeKind
pub fn control_dependency() -> EdgeKind
pub fn ownership_transfer_dependency() -> EdgeKind
pub fn synchronization_dependency() -> EdgeKind

工作流

Workflow

带标签和执行策略的任务图。

pub struct Workflow {
  label : String
  policy : @shared.ExecutionPolicy
  capabilities : Array[Capability]
  nodes : Array[Node]
  edges : Array[Edge]
} derive(Eq, @debug.Debug)
pub fn Workflow::new(String, policy? : @shared.ExecutionPolicy) -> Self

Workflow::new 创建一个空图;policy 默认为 @shared.native_policy()。

Workflow::add_capability、Workflow::add_node 和 Workflow::add_edge

追加一个能力、节点或边,并返回同一个工作流。

pub fn Workflow::add_capability(Self, Capability) -> Self
pub fn Workflow::add_node(Self, Node) -> Self
pub fn Workflow::add_edge(Self, Edge) -> Self

这些方法就地向工作流的数组追加元素:接收者与结果是同一个值,对该工作流的所有其他引用都能看到新增的内容。它们不做任何检查。

Workflow::label、Workflow::policy、Workflow::capabilities、Workflow::nodes 和 Workflow::edges

返回工作流的各个字段。返回的数组就是工作流自己的数组,而不是副本。

pub fn Workflow::label(Self) -> String
pub fn Workflow::policy(Self) -> @shared.ExecutionPolicy
pub fn Workflow::capabilities(Self) -> Array[Capability]
pub fn Workflow::nodes(Self) -> Array[Node]
pub fn Workflow::edges(Self) -> Array[Edge]
test "building a workflow" {
  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()))
  assert_eq(graph.nodes().length(), 2)
  assert_eq(graph.label(), "fork-join")
}

验证

WorkflowIssue

在工作流中发现的一个问题。

pub enum WorkflowIssue {
  EmptyWorkflow
  InvalidNodeReference(node_id~ : Int)
  SelfEdge(node_id~ : Int)
  CyclicDependency
  DuplicateNodeId(node_id~ : Int)
  DuplicateCapabilityId(capability_id~ : Int)
  MissingCapability(capability_id~ : Int)
  MissingNodeCapability(node_id~ : Int)
  InvalidDependency(from_id~ : Int, to_id~ : Int)
  UnsupportedSharedWrite(node_id~ : Int)
  InvalidCapabilityAccess(capability_id~ : Int)
  InvalidChannelNode(node_id~ : Int)
  InvalidMutexNode(node_id~ : Int)
  InvalidCondvarNode(node_id~ : Int)
  InvalidBarrierNode(node_id~ : Int)
  UnsupportedCapability(kind~ : CapabilityKind)
  InvalidPolicy(issue~ : @shared.PolicyIssue)
  InvalidComputePlan(node_id~ : Int, issue~ : @plan.ValidationIssue)
} derive(Eq, @debug.Debug)

validate

返回工作流的所有问题,若没有则返回空数组。

pub fn validate(Workflow) -> Array[WorkflowIssue]

各项检查,按其问题出现的顺序:

  1. 没有节点时报告 EmptyWorkflow,然后对工作流策略运行 @shared.validate_policy,每个问题报告一个 InvalidPolicy。
  2. 对每个能力:Opaque 报告 UnsupportedCapability;访问模式对该种类无效时报告 InvalidCapabilityAccess;其 id 出现不止一次时报告 DuplicateCapabilityId(每次出现报告一次)。
  3. 对每个节点:DuplicateNodeId(每次出现一次);种类需要能力却没有时报告 MissingNodeCapability;指定的能力不存在时报告 MissingCapability;能力种类不对时报告 InvalidChannelNode、InvalidMutexNode、InvalidCondvarNode 或 InvalidBarrierNode(其他节点种类则报告端点相同的 InvalidDependency);WriteShared 节点指定了 ReadOnly 能力时报告 UnsupportedSharedWrite;对计算节点的计划运行 @plan.validate,每个问题报告一个 InvalidComputePlan。
  4. 对每条边:from == to 时报告 SelfEdge;每个不是节点的端点报告一个 InvalidNodeReference;任一端点不是节点时报告 InvalidDependency。
  5. 图中没有入边为零的节点,或者同一对节点之间有两条方向相反的边时,报告 CyclicDependency。

第 5 步的环检查并不能找出所有的环:从某个源节点可达的、长度为三或以上的环会通过检查。原生运行时会完整地检查环,并在提交时拒绝这样的图。工作流设计讨论了这两种检查。

is_ready

当 validate 没有返回任何问题时返回 true。

pub fn is_ready(Workflow) -> Bool

is_invalid_compute_plan、is_cyclic_dependency、is_missing_capability 和 is_unsupported_shared_write

判断一个 WorkflowIssue 是哪种问题。

pub fn is_invalid_compute_plan(WorkflowIssue) -> Bool
pub fn is_cyclic_dependency(WorkflowIssue) -> Bool
pub fn is_missing_capability(WorkflowIssue) -> Bool
pub fn is_unsupported_shared_write(WorkflowIssue) -> Bool

is_missing_capability 同时匹配 MissingCapability 和 MissingNodeCapability。

test "validation" {
  let graph = @workflow.Workflow::new("bad").add_node(
    @workflow.write_shared_node(1, "write", capability=7),
  )
  let issues = @workflow.validate(graph)
  debug_inspect(issues, content="[MissingCapability(capability_id=7)]")
  assert_true(issues.any(@workflow.is_missing_capability))
  assert_true(!@workflow.is_ready(graph))
}

提交

Submission

提交的记录:工作流、目标后端、是否被接受、是否已完成,以及发现的问题。

pub struct Submission {
  workflow : Workflow
  backend : @shared.BackendTarget
  accepted : Bool
  completed : Bool
  issues : Array[WorkflowIssue]
} derive(Eq, @debug.Debug)
pub fn Submission::accepted(Self) -> Bool
pub fn Submission::backend(Self) -> @shared.BackendTarget
pub fn Submission::completed(Self) -> Bool
pub fn Submission::issues(Self) -> Array[WorkflowIssue]

submit

验证工作流,并记录后端是否接受它。

pub fn submit(Workflow, backend? : @shared.BackendTarget) -> Submission

当且仅当 validate 没有返回问题且后端(默认为 Native)是 Native 时,提交被接受。completed 总是 false:submit 不会运行工作流。

test "submit" {
  let graph = @workflow.Workflow::new("one").add_node(
    @workflow.spawn_node(1, "spawn"),
  )
  let native = @workflow.submit(graph)
  assert_true(native.accepted() && !native.completed())
  let js = @workflow.submit(graph, backend=@shared.javascript_target())
  assert_true(!js.accepted())
  assert_eq(js.issues().length(), 0)
}

相等性与包信息

T::equal

结构相等,作为方法提升到本包的每个类型上:AccessMode、Capability、CapabilityKind、Edge、EdgeKind、Node、NodeKind、Submission、Workflow 和 WorkflowIssue。请使用 == 和 !=。

pub fn Workflow::equal(Self, Self) -> Bool

package_name

返回 "workflow"。

pub fn package_name() -> String