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)
4 種類のエッジはどれも端点を同じように順序付けます。種類は理由を記録するもので、そのままランタイムに渡されます。
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
ワークフローに見つかった問題の 1 つです。
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を 1 つ。 - 各ケイパビリティについて、
OpaqueならUnsupportedCapability、アクセスモードがその種類に有効でなければInvalidCapabilityAccess、id が 2 回以上現れればDuplicateCapabilityId(出現ごとに 1 回)。 - 各ノードについて、
DuplicateNodeId(出現ごとに 1 回)、種類がケイパビリティを必要とするのに持たなければMissingNodeCapability、指定したケイパビリティが存在しなければMissingCapability、ケイパビリティの種類が合わなければInvalidChannelNode、InvalidMutexNode、InvalidCondvarNode、InvalidBarrierNodeのいずれか(ほかのノード種類では端点が等しいInvalidDependency)、WriteSharedノードがReadOnlyのケイパビリティを指定していればUnsupportedSharedWrite、そして計算ノードのプランに対する@plan.validateの問題ごとにInvalidComputePlanを 1 つ。 - 各エッジについて、
from == toならSelfEdge、ノードでない端点ごとにInvalidNodeReference、どちらかがノードでなければInvalidDependency。 - 入ってくるエッジのないノードがグラフに 1 つもないとき、または同じ 2 ノードの間に逆向きのエッジが 2 本あるとき
CyclicDependency。
手順 5 の閉路検査はすべての閉路を見つけるわけではありません。ソースから到達できる長さ 3 以上の閉路は通ってしまいます。ネイティブランタイムは閉路を完全に検査し、投入時にそのようなグラフを拒否します。両方の検査はワークフローの設計で説明しています。
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