///|
/// A task definition in workflow-spec-v1 input JSON.
///
/// Unlike `TaskNode`, this type models user-supplied workflow structure and does
/// not carry runtime execution status.
pub struct WorkflowTaskSpec {
priv id : TaskId
priv title : String
priv description : String
priv inputs : Array[String]
priv outputs : Array[String]
priv tags : Array[String]
} derive(Eq, Debug)
///|
/// Workflow input definition, separated from run-snapshot JSON output.
pub struct WorkflowSpec {
priv tasks : Array[WorkflowTaskSpec]
priv dependencies : Array[Dependency]
} derive(Eq, Debug)
///|
/// Errors raised while parsing and validating workflow-spec-v1 JSON.
pub(all) enum WorkflowSpecError {
InvalidWorkflowJson(String)
UnsupportedWorkflowSchema(String)
MissingWorkflowField(String)
InvalidWorkflowField(String)
WorkflowGraphError(GraphError)
} derive(Eq, Debug)
///|
/// Build a workflow task spec.
pub fn WorkflowTaskSpec::new(id : String, title : String) -> WorkflowTaskSpec {
{
id: TaskId::new(id),
title,
description: "",
inputs: [],
outputs: [],
tags: [],
}
}
///|
/// Add a short description to a workflow task spec.
pub fn WorkflowTaskSpec::with_description(
self : WorkflowTaskSpec,
description : String,
) -> WorkflowTaskSpec {
{ ..self, description, }
}
///|
/// Add input labels or artifacts to a workflow task spec.
pub fn WorkflowTaskSpec::with_inputs(
self : WorkflowTaskSpec,
inputs : Array[String],
) -> WorkflowTaskSpec {
{ ..self, inputs: inputs.copy() }
}
///|
/// Add output labels or artifacts to a workflow task spec.
pub fn WorkflowTaskSpec::with_outputs(
self : WorkflowTaskSpec,
outputs : Array[String],
) -> WorkflowTaskSpec {
{ ..self, outputs: outputs.copy() }
}
///|
/// Add topic tags to a workflow task spec.
pub fn WorkflowTaskSpec::with_tags(
self : WorkflowTaskSpec,
tags : Array[String],
) -> WorkflowTaskSpec {
{ ..self, tags: tags.copy() }
}
///|
/// Return this workflow task id.
pub fn WorkflowTaskSpec::id(self : WorkflowTaskSpec) -> TaskId {
self.id
}
///|
/// Return this workflow task title.
pub fn WorkflowTaskSpec::title(self : WorkflowTaskSpec) -> String {
self.title
}
///|
/// Return this workflow task description.
pub fn WorkflowTaskSpec::description(self : WorkflowTaskSpec) -> String {
self.description
}
///|
/// Return a detached copy of this workflow task's input labels.
pub fn WorkflowTaskSpec::inputs(self : WorkflowTaskSpec) -> Array[String] {
self.inputs.copy()
}
///|
/// Return a detached copy of this workflow task's output labels.
pub fn WorkflowTaskSpec::outputs(self : WorkflowTaskSpec) -> Array[String] {
self.outputs.copy()
}
///|
/// Return a detached copy of this workflow task's tags.
pub fn WorkflowTaskSpec::tags(self : WorkflowTaskSpec) -> Array[String] {
self.tags.copy()
}
///|
/// Return a detached copy of a workflow task spec.
pub fn WorkflowTaskSpec::snapshot(self : WorkflowTaskSpec) -> WorkflowTaskSpec {
{
..self,
inputs: self.inputs.copy(),
outputs: self.outputs.copy(),
tags: self.tags.copy(),
}
}
///|
/// Build a workflow spec from tasks and dependency edges.
pub fn WorkflowSpec::new(
tasks : Array[WorkflowTaskSpec],
dependencies : Array[Dependency],
) -> WorkflowSpec {
let task_copies : Array[WorkflowTaskSpec] = []
for task in tasks {
task_copies.push(task.snapshot())
}
{ tasks: task_copies, dependencies: dependencies.copy() }
}
///|
/// Return a detached copy of workflow task specs.
pub fn WorkflowSpec::tasks(self : WorkflowSpec) -> Array[WorkflowTaskSpec] {
let out : Array[WorkflowTaskSpec] = []
for task in self.tasks {
out.push(task.snapshot())
}
out
}
///|
/// Return a detached copy of workflow dependency edges.
pub fn WorkflowSpec::dependencies(self : WorkflowSpec) -> Array[Dependency] {
self.dependencies.copy()
}
///|
/// Parse, type-check, and validate a workflow-spec-v1 JSON string.
pub fn WorkflowSpec::from_json(
input : String,
) -> Result[WorkflowSpec, WorkflowSpecError] {
let json = @json.parse(input) catch {
err => return Err(InvalidWorkflowJson("\{err}"))
}
let spec = match parse_workflow_spec_json(json) {
Err(err) => return Err(err)
Ok(spec) => spec
}
match spec.to_graph() {
Err(err) => Err(WorkflowGraphError(err))
Ok(_) => Ok(spec)
}
}
///|
/// Render this workflow spec as workflow-spec-v1 JSON.
pub fn WorkflowSpec::to_json(self : WorkflowSpec) -> String {
let out = StringBuilder()
out <+ "{\n"
out <+ " \"schema_version\": 1,\n"
out <+ " \"tasks\": ["
for i = 0; i < self.tasks.length(); i = i + 1 {
if i > 0 {
out <+ ", "
}
let task = self.tasks[i]
let id = escape_json(task.id.value)
let title = escape_json(task.title)
let description = escape_json(task.description)
out <+ "{"
out <+ "\"id\":\"\{id}\","
out <+ "\"title\":\"\{title}\","
out <+ "\"description\":\"\{description}\","
out <+ "\"inputs\":"
write_json_strings(out, task.inputs)
out <+ ",\"outputs\":"
write_json_strings(out, task.outputs)
out <+ ",\"tags\":"
write_json_strings(out, task.tags)
out <+ "}"
}
out <+ "],\n"
out <+ " \"dependencies\": ["
for i = 0; i < self.dependencies.length(); i = i + 1 {
if i > 0 {
out <+ ", "
}
let dep = self.dependencies[i]
let before = escape_json(dep.before.value)
let after = escape_json(dep.after.value)
out <+ "{\"before\":\"\{before}\",\"after\":\"\{after}\"}"
}
out <+ "]\n"
out <+ "}\n"
out.to_string()
}
///|
/// Convert this workflow definition into a validated flow graph.
pub fn WorkflowSpec::to_graph(
self : WorkflowSpec,
) -> Result[FlowGraph, GraphError] {
let graph = FlowGraph::new()
for task in self.tasks {
match
graph.add_task(
TaskNode::new(task.id.value, task.title)
.with_description(task.description)
.with_inputs(task.inputs)
.with_outputs(task.outputs)
.with_tags(task.tags),
) {
Err(err) => return Err(err)
Ok(_) => ()
}
}
for dep in self.dependencies {
match graph.add_dependency(dep.before, dep.after) {
Err(err) => return Err(err)
Ok(_) => ()
}
}
match graph.validate() {
Err(err) => Err(err)
Ok(_) => Ok(graph)
}
}
///|
/// Convert a graph into workflow-spec-v1 input JSON semantics.
///
/// Runtime task status is intentionally not preserved in workflow specs.
pub fn FlowGraph::to_workflow_spec(self : FlowGraph) -> WorkflowSpec {
let tasks : Array[WorkflowTaskSpec] = []
for task in self.tasks {
tasks.push(
WorkflowTaskSpec::new(task.id.value, task.title)
.with_description(task.description)
.with_inputs(task.inputs)
.with_outputs(task.outputs)
.with_tags(task.tags),
)
}
WorkflowSpec::new(tasks, self.dependencies)
}
///|
/// Parse workflow-spec-v1 JSON and return a validated graph.
pub fn FlowGraph::from_workflow_json(
input : String,
) -> Result[FlowGraph, WorkflowSpecError] {
match WorkflowSpec::from_json(input) {
Err(err) => Err(err)
Ok(spec) =>
match spec.to_graph() {
Err(err) => Err(WorkflowGraphError(err))
Ok(graph) => Ok(graph)
}
}
}
///|
/// Return a readable diagnostic for workflow spec errors.
pub fn WorkflowSpecError::message(self : WorkflowSpecError) -> String {
match self {
InvalidWorkflowJson(message) => "invalid workflow JSON: \{message}"
UnsupportedWorkflowSchema(version) =>
"unsupported workflow schema_version: \{version}"
MissingWorkflowField(path) => "missing workflow field: \{path}"
InvalidWorkflowField(path) => "invalid workflow field: \{path}"
WorkflowGraphError(err) => err.message()
}
}
///|
fn parse_workflow_spec_json(
json : Json,
) -> Result[WorkflowSpec, WorkflowSpecError] {
guard json is Object(obj) else {
return Err(InvalidWorkflowField("$: expected object"))
}
match obj.get("schema_version") {
None => return Err(MissingWorkflowField("schema_version"))
Some(Number(version, ..)) =>
if version != 1.0 {
return Err(UnsupportedWorkflowSchema("\{version}"))
}
Some(_) =>
return Err(InvalidWorkflowField("schema_version: expected number 1"))
}
let tasks = match obj.get("tasks") {
None => return Err(MissingWorkflowField("tasks"))
Some(Array(items)) => parse_workflow_tasks(items)
Some(_) => return Err(InvalidWorkflowField("tasks: expected array"))
}
let dependencies = match obj.get("dependencies") {
None => return Err(MissingWorkflowField("dependencies"))
Some(Array(items)) => parse_workflow_dependencies(items)
Some(_) => return Err(InvalidWorkflowField("dependencies: expected array"))
}
match tasks {
Err(err) => Err(err)
Ok(tasks) =>
match dependencies {
Err(err) => Err(err)
Ok(dependencies) => Ok(WorkflowSpec::new(tasks, dependencies))
}
}
}
///|
fn parse_workflow_tasks(
items : Array[Json],
) -> Result[Array[WorkflowTaskSpec], WorkflowSpecError] {
let tasks : Array[WorkflowTaskSpec] = []
for i = 0; i < items.length(); i = i + 1 {
guard items[i] is Object(obj) else {
return Err(InvalidWorkflowField("tasks[\{i}]: expected object"))
}
let id = match required_string(obj, "tasks[\{i}].id") {
Err(err) => return Err(err)
Ok(value) => value
}
let title = match required_string(obj, "tasks[\{i}].title") {
Err(err) => return Err(err)
Ok(value) => value
}
let description = match
optional_string(obj, "description", "tasks[\{i}].description", "") {
Err(err) => return Err(err)
Ok(value) => value
}
let inputs = match
optional_string_array(obj, "inputs", "tasks[\{i}].inputs") {
Err(err) => return Err(err)
Ok(values) => values
}
let outputs = match
optional_string_array(obj, "outputs", "tasks[\{i}].outputs") {
Err(err) => return Err(err)
Ok(values) => values
}
let tags = match optional_string_array(obj, "tags", "tasks[\{i}].tags") {
Err(err) => return Err(err)
Ok(values) => values
}
tasks.push(
WorkflowTaskSpec::new(id, title)
.with_description(description)
.with_inputs(inputs)
.with_outputs(outputs)
.with_tags(tags),
)
}
Ok(tasks)
}
///|
fn parse_workflow_dependencies(
items : Array[Json],
) -> Result[Array[Dependency], WorkflowSpecError] {
let dependencies : Array[Dependency] = []
for i = 0; i < items.length(); i = i + 1 {
guard items[i] is Object(obj) else {
return Err(InvalidWorkflowField("dependencies[\{i}]: expected object"))
}
let before = match required_string(obj, "dependencies[\{i}].before") {
Err(err) => return Err(err)
Ok(value) => value
}
let after = match required_string(obj, "dependencies[\{i}].after") {
Err(err) => return Err(err)
Ok(value) => value
}
dependencies.push(Dependency::new(TaskId::new(before), TaskId::new(after)))
}
Ok(dependencies)
}
///|
fn required_string(
obj : Map[String, Json],
path : String,
) -> Result[String, WorkflowSpecError] {
let field = last_path_component(path)
match obj.get(field) {
None => Err(MissingWorkflowField(path))
Some(String(value)) => Ok(value)
Some(_) => Err(InvalidWorkflowField("\{path}: expected string"))
}
}
///|
fn optional_string(
obj : Map[String, Json],
field : String,
path : String,
default : String,
) -> Result[String, WorkflowSpecError] {
match obj.get(field) {
None => Ok(default)
Some(String(value)) => Ok(value)
Some(_) => Err(InvalidWorkflowField("\{path}: expected string"))
}
}
///|
fn optional_string_array(
obj : Map[String, Json],
field : String,
path : String,
) -> Result[Array[String], WorkflowSpecError] {
match obj.get(field) {
None => Ok([])
Some(value) => parse_string_array(value, path)
}
}
///|
fn parse_string_array(
value : Json,
path : String,
) -> Result[Array[String], WorkflowSpecError] {
guard value is Array(items) else {
return Err(InvalidWorkflowField("\{path}: expected string array"))
}
let out : Array[String] = []
for i = 0; i < items.length(); i = i + 1 {
match items[i] {
String(value) => out.push(value)
_ => return Err(InvalidWorkflowField("\{path}[\{i}]: expected string"))
}
}
Ok(out)
}
///|
fn last_path_component(path : String) -> String {
let mut start = 0
for i = 0; i < path.length(); i = i + 1 {
if path[i] == '.' {
start = i + 1
}
}
path[start:].to_owned()
}