///|
/// A validated run-snapshot-v1 artifact containing graph, plan, and trace data.
pub struct RunSnapshot {
priv graph : FlowGraph
priv plan : ExecutionPlan
priv trace : Trace
} derive(Debug)
///|
/// Errors raised while parsing or auditing run-snapshot-v1 JSON.
pub(all) enum RunSnapshotError {
InvalidRunSnapshotJson(String)
UnsupportedRunSnapshotSchema(String)
MissingRunSnapshotField(String)
InvalidRunSnapshotField(String)
RunSnapshotGraphError(GraphError)
RunSnapshotPlanMismatch
RunSnapshotUnknownTraceTask(TaskId)
InvalidTraceLifecycle(TaskId, TraceEventType)
RunSnapshotStatusMismatch(TaskId, TaskStatus, String)
} derive(Eq, Debug)
///|
/// Parse and audit a run-snapshot-v1 JSON string.
pub fn RunSnapshot::from_json(
input : String,
) -> Result[RunSnapshot, RunSnapshotError] {
let json = @json.parse(input) catch {
err => return Err(InvalidRunSnapshotJson("\{err}"))
}
let snapshot = match parse_run_snapshot_json(json) {
Err(err) => return Err(err)
Ok(snapshot) => snapshot
}
match snapshot.audit() {
Err(err) => Err(err)
Ok(_) => Ok(snapshot)
}
}
///|
/// Return a detached graph copy.
pub fn RunSnapshot::graph(self : RunSnapshot) -> FlowGraph {
self.graph.snapshot()
}
///|
/// Return a detached execution-plan copy.
pub fn RunSnapshot::plan(self : RunSnapshot) -> ExecutionPlan {
self.plan.snapshot()
}
///|
/// Return a detached trace copy.
pub fn RunSnapshot::trace(self : RunSnapshot) -> Trace {
self.trace.snapshot()
}
///|
/// Render this audited snapshot in canonical run-snapshot-v1 form.
pub fn RunSnapshot::to_json(self : RunSnapshot) -> String {
self.graph.render_json(self.plan, self.trace, versioned=true)
}
///|
/// Validate graph structure, serialized plan, trace references, and lifecycle.
pub fn RunSnapshot::audit(self : RunSnapshot) -> Result[Unit, RunSnapshotError] {
match self.graph.validate() {
Err(err) => return Err(RunSnapshotGraphError(err))
Ok(_) => ()
}
match self.graph.validate_plan(self.plan) {
Err(_) => return Err(RunSnapshotPlanMismatch)
Ok(_) => ()
}
match self.graph.validate_trace(self.trace) {
Err(UnknownTraceTask(id)) => return Err(RunSnapshotUnknownTraceTask(id))
Err(SnapshotGraphError(err)) => return Err(RunSnapshotGraphError(err))
Err(_) => return Err(RunSnapshotPlanMismatch)
Ok(_) => ()
}
audit_trace_lifecycle(self.graph, self.trace)
}
///|
/// Return a readable diagnostic for snapshot import and audit errors.
pub fn RunSnapshotError::message(self : RunSnapshotError) -> String {
match self {
InvalidRunSnapshotJson(message) => "invalid run snapshot JSON: \{message}"
UnsupportedRunSnapshotSchema(version) =>
"unsupported run snapshot schema_version: \{version}"
MissingRunSnapshotField(path) => "missing run snapshot field: \{path}"
InvalidRunSnapshotField(path) => "invalid run snapshot field: \{path}"
RunSnapshotGraphError(err) => err.message()
RunSnapshotPlanMismatch =>
"serialized execution plan does not match the snapshot graph"
RunSnapshotUnknownTraceTask(id) =>
"snapshot trace references unknown task: \{id.value()}"
InvalidTraceLifecycle(id, event_type) =>
"invalid trace lifecycle for \{id.value()}: \{event_type.label()}"
RunSnapshotStatusMismatch(id, status, expected) =>
"snapshot status mismatch for \{id.value()}: \{status.label()}, expected \{expected}"
}
}
///|
fn parse_run_snapshot_json(
json : Json,
) -> Result[RunSnapshot, RunSnapshotError] {
guard json is Object(obj) else {
return Err(InvalidRunSnapshotField("$: expected object"))
}
match
snapshot_reject_unknown_fields(obj, "$", [
"schema_version", "order", "batches", "tasks", "dependencies", "trace",
]) {
Err(err) => return Err(err)
Ok(_) => ()
}
match obj.get("schema_version") {
None => return Err(MissingRunSnapshotField("schema_version"))
Some(Number(version, ..)) =>
if version != 1.0 {
return Err(UnsupportedRunSnapshotSchema("\{version}"))
}
Some(_) =>
return Err(InvalidRunSnapshotField("schema_version: expected number 1"))
}
let order = match snapshot_string_array_field(obj, "order") {
Err(err) => return Err(err)
Ok(value) => ids_from_strings(value)
}
let batches = match snapshot_batches_field(obj, "batches") {
Err(err) => return Err(err)
Ok(value) => value
}
let tasks = match snapshot_tasks_field(obj, "tasks") {
Err(err) => return Err(err)
Ok(value) => value
}
let dependencies = match snapshot_dependencies_field(obj, "dependencies") {
Err(err) => return Err(err)
Ok(value) => value
}
let events = match snapshot_trace_field(obj, "trace") {
Err(err) => return Err(err)
Ok(value) => value
}
let graph = FlowGraph::new()
for task in tasks {
match graph.add_task(task) {
Err(err) => return Err(RunSnapshotGraphError(err))
Ok(_) => ()
}
}
for dependency in dependencies {
match graph.add_dependency(dependency.before, dependency.after) {
Err(err) => return Err(RunSnapshotGraphError(err))
Ok(_) => ()
}
}
let trace = Trace::new()
for event in events {
trace.record(event)
}
Ok({ graph, plan: { order, batches }, trace })
}
///|
fn snapshot_tasks_field(
obj : Map[String, Json],
field : String,
) -> Result[Array[TaskNode], RunSnapshotError] {
let items = match obj.get(field) {
None => return Err(MissingRunSnapshotField(field))
Some(Array(items)) => items
Some(_) => return Err(InvalidRunSnapshotField("\{field}: expected array"))
}
let tasks : Array[TaskNode] = []
for i = 0; i < items.length(); i = i + 1 {
guard items[i] is Object(task_obj) else {
return Err(InvalidRunSnapshotField("tasks[\{i}]: expected object"))
}
match
snapshot_reject_unknown_fields(task_obj, "tasks[\{i}]", [
"id", "title", "description", "status", "status_kind", "status_reason", "inputs",
"outputs", "tags",
]) {
Err(err) => return Err(err)
Ok(_) => ()
}
let id = match snapshot_required_string(task_obj, "id", "tasks[\{i}].id") {
Err(err) => return Err(err)
Ok(value) => value
}
let title = match
snapshot_required_string(task_obj, "title", "tasks[\{i}].title") {
Err(err) => return Err(err)
Ok(value) => value
}
let description = match
snapshot_required_string(
task_obj,
"description",
"tasks[\{i}].description",
) {
Err(err) => return Err(err)
Ok(value) => value
}
let display_status = match
snapshot_required_string(task_obj, "status", "tasks[\{i}].status") {
Err(err) => return Err(err)
Ok(value) => value
}
let status_kind = match
snapshot_required_string(
task_obj,
"status_kind",
"tasks[\{i}].status_kind",
) {
Err(err) => return Err(err)
Ok(value) => value
}
let status_reason = match
snapshot_nullable_string(
task_obj,
"status_reason",
"tasks[\{i}].status_reason",
) {
Err(err) => return Err(err)
Ok(value) => value
}
let status = match
parse_snapshot_status(status_kind, status_reason, "tasks[\{i}]") {
Err(err) => return Err(err)
Ok(value) => value
}
if status.label() != display_status {
return Err(
InvalidRunSnapshotField("tasks[\{i}].status: inconsistent status label"),
)
}
let inputs = match
snapshot_string_array_field(task_obj, "inputs", path="tasks[\{i}].inputs") {
Err(err) => return Err(err)
Ok(value) => value
}
let outputs = match
snapshot_string_array_field(
task_obj,
"outputs",
path="tasks[\{i}].outputs",
) {
Err(err) => return Err(err)
Ok(value) => value
}
let tags = match
snapshot_string_array_field(task_obj, "tags", path="tasks[\{i}].tags") {
Err(err) => return Err(err)
Ok(value) => value
}
tasks.push(
TaskNode::new(id, title)
.with_description(description)
.with_inputs(inputs)
.with_outputs(outputs)
.with_tags(tags)
.with_status(status),
)
}
Ok(tasks)
}
///|
fn snapshot_dependencies_field(
obj : Map[String, Json],
field : String,
) -> Result[Array[Dependency], RunSnapshotError] {
let items = match obj.get(field) {
None => return Err(MissingRunSnapshotField(field))
Some(Array(items)) => items
Some(_) => return Err(InvalidRunSnapshotField("\{field}: expected array"))
}
let dependencies : Array[Dependency] = []
for i = 0; i < items.length(); i = i + 1 {
guard items[i] is Object(dep_obj) else {
return Err(InvalidRunSnapshotField("dependencies[\{i}]: expected object"))
}
match
snapshot_reject_unknown_fields(dep_obj, "dependencies[\{i}]", [
"before", "after",
]) {
Err(err) => return Err(err)
Ok(_) => ()
}
let before = match
snapshot_required_string(dep_obj, "before", "dependencies[\{i}].before") {
Err(err) => return Err(err)
Ok(value) => value
}
let after = match
snapshot_required_string(dep_obj, "after", "dependencies[\{i}].after") {
Err(err) => return Err(err)
Ok(value) => value
}
dependencies.push(Dependency::new(TaskId::new(before), TaskId::new(after)))
}
Ok(dependencies)
}
///|
fn snapshot_trace_field(
obj : Map[String, Json],
field : String,
) -> Result[Array[TraceEvent], RunSnapshotError] {
let items = match obj.get(field) {
None => return Err(MissingRunSnapshotField(field))
Some(Array(items)) => items
Some(_) => return Err(InvalidRunSnapshotField("\{field}: expected array"))
}
let events : Array[TraceEvent] = []
for i = 0; i < items.length(); i = i + 1 {
guard items[i] is Object(event_obj) else {
return Err(InvalidRunSnapshotField("trace[\{i}]: expected object"))
}
match
snapshot_reject_unknown_fields(event_obj, "trace[\{i}]", [
"task_id", "event_type", "timestamp", "message",
]) {
Err(err) => return Err(err)
Ok(_) => ()
}
let task_id = match
snapshot_required_string(event_obj, "task_id", "trace[\{i}].task_id") {
Err(err) => return Err(err)
Ok(value) => value
}
let event_label = match
snapshot_required_string(
event_obj,
"event_type",
"trace[\{i}].event_type",
) {
Err(err) => return Err(err)
Ok(value) => value
}
let event_type = match
parse_snapshot_event_type(event_label, "trace[\{i}]") {
Err(err) => return Err(err)
Ok(value) => value
}
let timestamp = match
snapshot_required_string(event_obj, "timestamp", "trace[\{i}].timestamp") {
Err(err) => return Err(err)
Ok(value) => value
}
let message = match
snapshot_required_string(event_obj, "message", "trace[\{i}].message") {
Err(err) => return Err(err)
Ok(value) => value
}
events.push(
TraceEvent::new(TaskId::new(task_id), event_type, message, timestamp),
)
}
Ok(events)
}
///|
fn snapshot_batches_field(
obj : Map[String, Json],
field : String,
) -> Result[Array[Array[TaskId]], RunSnapshotError] {
let items = match obj.get(field) {
None => return Err(MissingRunSnapshotField(field))
Some(Array(items)) => items
Some(_) => return Err(InvalidRunSnapshotField("\{field}: expected array"))
}
let batches : Array[Array[TaskId]] = []
for i = 0; i < items.length(); i = i + 1 {
let values = match snapshot_parse_string_array(items[i], "batches[\{i}]") {
Err(err) => return Err(err)
Ok(value) => value
}
batches.push(ids_from_strings(values))
}
Ok(batches)
}
///|
fn snapshot_string_array_field(
obj : Map[String, Json],
field : String,
path? : String = field,
) -> Result[Array[String], RunSnapshotError] {
match obj.get(field) {
None => Err(MissingRunSnapshotField(path))
Some(value) => snapshot_parse_string_array(value, path)
}
}
///|
fn snapshot_parse_string_array(
value : Json,
path : String,
) -> Result[Array[String], RunSnapshotError] {
guard value is Array(items) else {
return Err(InvalidRunSnapshotField("\{path}: expected string array"))
}
let values : Array[String] = []
for i = 0; i < items.length(); i = i + 1 {
match items[i] {
String(value) => values.push(value)
_ => return Err(InvalidRunSnapshotField("\{path}[\{i}]: expected string"))
}
}
Ok(values)
}
///|
fn snapshot_required_string(
obj : Map[String, Json],
field : String,
path : String,
) -> Result[String, RunSnapshotError] {
match obj.get(field) {
None => Err(MissingRunSnapshotField(path))
Some(String(value)) => Ok(value)
Some(_) => Err(InvalidRunSnapshotField("\{path}: expected string"))
}
}
///|
fn snapshot_reject_unknown_fields(
obj : Map[String, Json],
path : String,
allowed : Array[String],
) -> Result[Unit, RunSnapshotError] {
for key in obj.keys() {
let mut known = false
for candidate in allowed {
if key == candidate {
known = true
break
}
}
if !known {
return Err(InvalidRunSnapshotField("\{path}.\{key}: unknown field"))
}
}
Ok(())
}
///|
fn snapshot_nullable_string(
obj : Map[String, Json],
field : String,
path : String,
) -> Result[String?, RunSnapshotError] {
match obj.get(field) {
None => Err(MissingRunSnapshotField(path))
Some(Null) => Ok(None)
Some(String(value)) => Ok(Some(value))
Some(_) => Err(InvalidRunSnapshotField("\{path}: expected string or null"))
}
}
///|
fn parse_snapshot_status(
kind : String,
reason : String?,
path : String,
) -> Result[TaskStatus, RunSnapshotError] {
match (kind, reason) {
("pending", None) => Ok(Pending)
("ready", None) => Ok(Ready)
("running", None) => Ok(Running)
("succeeded", None) => Ok(Succeeded)
("failed", Some(reason)) => Ok(Failed(reason))
("skipped", Some(reason)) => Ok(Skipped(reason))
("failed" | "skipped", None) =>
Err(InvalidRunSnapshotField("\{path}.status_reason: expected string"))
("pending" | "ready" | "running" | "succeeded", Some(_)) =>
Err(InvalidRunSnapshotField("\{path}.status_reason: expected null"))
_ => Err(InvalidRunSnapshotField("\{path}.status_kind: unknown status"))
}
}
///|
fn parse_snapshot_event_type(
label : String,
path : String,
) -> Result[TraceEventType, RunSnapshotError] {
match label {
"planned" => Ok(Planned)
"started" => Ok(Started)
"completed" => Ok(Completed)
"failed" => Ok(FailedEvent)
"skipped" => Ok(SkippedEvent)
"note" => Ok(Note)
_ => Err(InvalidRunSnapshotField("\{path}.event_type: unknown event"))
}
}
///|
fn ids_from_strings(values : Array[String]) -> Array[TaskId] {
let ids : Array[TaskId] = []
for value in values {
ids.push(TaskId::new(value))
}
ids
}
///|
fn audit_trace_lifecycle(
graph : FlowGraph,
trace : Trace,
) -> Result[Unit, RunSnapshotError] {
for task in graph.tasks {
let mut state = "none"
for event in trace.events {
if event.task_id != task.id || event.event_type == Note {
continue
}
let next = match (state, event.event_type) {
("none", Planned) => "planned"
("none" | "planned" | "failed", Started) => "running"
("none" | "planned", SkippedEvent) => "skipped"
("running", Completed) => "succeeded"
("running", FailedEvent) => "failed"
_ => return Err(InvalidTraceLifecycle(task.id, event.event_type))
}
state = next
}
let matches = match (state, task.status) {
("none" | "planned", Pending | Ready) => true
("running", Running) => true
("succeeded", Succeeded) => true
("failed", Failed(_)) => true
("skipped", Skipped(_)) => true
_ => false
}
if !matches {
return Err(RunSnapshotStatusMismatch(task.id, task.status, state))
}
}
Ok(())
}