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)
每种能力只允许部分访问模式:
| 种类 | 有效的访问模式 |
|---|---|
OwnedBuffer | ReadOnly, WriteOnly, ReadWrite |
SharedReadView | ReadOnly |
AtomicCell | ReadWrite |
Mutex, Condvar, RwLock, Semaphore, Barrier | SynchronizeOnly |
Channel | MoveOnly |
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, Recv | Channel |
Lock, Unlock | Mutex, RwLock |
Wait, Signal | Condvar |
Barrier | Barrier |
ReadShared | SharedReadView, AtomicCell |
WriteShared | AtomicCell、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]
各项检查,按其问题出现的顺序:
- 没有节点时报告
EmptyWorkflow,然后对工作流策略运行@shared.validate_policy,每个问题报告一个InvalidPolicy。 - 对每个能力:
Opaque报告UnsupportedCapability;访问模式对该种类无效时报告InvalidCapabilityAccess;其 id 出现不止一次时报告DuplicateCapabilityId(每次出现报告一次)。 - 对每个节点:
DuplicateNodeId(每次出现一次);种类需要能力却没有时报告MissingNodeCapability;指定的能力不存在时报告MissingCapability;能力种类不对时报告InvalidChannelNode、InvalidMutexNode、InvalidCondvarNode或InvalidBarrierNode(其他节点种类则报告端点相同的InvalidDependency);WriteShared节点指定了ReadOnly能力时报告UnsupportedSharedWrite;对计算节点的计划运行@plan.validate,每个问题报告一个InvalidComputePlan。 - 对每条边:
from == to时报告SelfEdge;每个不是节点的端点报告一个InvalidNodeReference;任一端点不是节点时报告InvalidDependency。 - 图中没有入边为零的节点,或者同一对节点之间有两条方向相反的边时,报告
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