use crate::eval::Operator;
use crate::model::{
is_safe_id, is_valid_vendor, CoreNodeType, FlowDefinition, FlowNode, FlowNodeType, SavedFlow,
SUPPORTED_SPEC_VERSIONS,
};
use crate::nodes::{BranchData, ConditionalData};
use std::collections::HashSet;
use std::fmt;
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum ValidationError {
InvalidFlowId(String),
UnsupportedSpecVersion(String),
MultipleEntryNodes(usize),
DuplicateNodeId(String),
DuplicateEdgeId(String),
DanglingEdgeSource {
edge: String,
source: String,
},
DanglingEdgeTarget {
edge: String,
target: String,
},
InvalidVendorNamespace {
node: String,
node_type: String,
},
V2NodeInV1Document {
node: String,
node_type: String,
},
InvalidNodeData {
node: String,
message: String,
},
UnknownOperator {
node: String,
operator: String,
},
MissingHandleEdge {
node: String,
handle: String,
},
InvalidRequires {
message: String,
},
InvalidSchedule {
message: String,
},
}
impl fmt::Display for ValidationError {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
match self {
ValidationError::InvalidFlowId(id) => {
write!(f, "invalid flow id {id:?}: must match [A-Za-z0-9-]{{1,64}}")
}
ValidationError::UnsupportedSpecVersion(v) => {
write!(f, "unsupported spec_version {v:?}; this parser supports {SUPPORTED_SPEC_VERSIONS:?}")
}
ValidationError::MultipleEntryNodes(n) => {
write!(f, "flow has {n} entry nodes; at most one is allowed")
}
ValidationError::DuplicateNodeId(id) => {
write!(f, "duplicate node id {id:?}")
}
ValidationError::DuplicateEdgeId(id) => {
write!(f, "duplicate edge id {id:?}")
}
ValidationError::DanglingEdgeSource { edge, source } => {
write!(f, "edge {edge:?} references unknown source node {source:?}")
}
ValidationError::DanglingEdgeTarget { edge, target } => {
write!(f, "edge {edge:?} references unknown target node {target:?}")
}
ValidationError::InvalidVendorNamespace { node, node_type } => {
write!(
f,
"node {node:?} has malformed custom node_type {node_type:?}: \
vendor prefix must match [a-z][a-z0-9_-]{{0,31}}"
)
}
ValidationError::V2NodeInV1Document { node, node_type } => {
write!(
f,
"node {node:?} uses v2 node type {node_type:?} but the document \
declares spec_version \"1\"; set spec_version to \"2\""
)
}
ValidationError::InvalidNodeData { node, message } => {
write!(f, "node {node:?} has invalid data: {message}")
}
ValidationError::UnknownOperator { node, operator } => {
write!(f, "node {node:?} uses unknown operator {operator:?}")
}
ValidationError::MissingHandleEdge { node, handle } => {
write!(f, "node {node:?} declares handle {handle:?} but no edge leaves it via that handle")
}
ValidationError::InvalidRequires { message } => {
write!(f, "invalid `requires`: {message}")
}
ValidationError::InvalidSchedule { message } => {
write!(f, "invalid schedule: {message}")
}
}
}
}
impl std::error::Error for ValidationError {}
pub fn validate(flow: &SavedFlow) -> Vec<ValidationError> {
let mut errors = Vec::new();
if !is_safe_id(&flow.id) {
errors.push(ValidationError::InvalidFlowId(flow.id.clone()));
}
if !SUPPORTED_SPEC_VERSIONS.contains(&flow.spec_version.as_str()) {
errors.push(ValidationError::UnsupportedSpecVersion(
flow.spec_version.clone(),
));
}
if let Some(req) = &flow.requires {
validate_requires(req, &mut errors);
}
validate_schedules(&flow.schedules, &mut errors);
validate_definition(&flow.flow, &flow.spec_version, &mut errors);
errors
}
fn validate_schedules(
schedules: &[crate::model::FlowScheduleSpec],
errors: &mut Vec<ValidationError>,
) {
use crate::model::ScheduleTrigger;
let mut seen: HashSet<&str> = HashSet::new();
for s in schedules {
if s.id.trim().is_empty() {
errors.push(ValidationError::InvalidSchedule {
message: "schedule id must not be empty".to_string(),
});
} else if !seen.insert(s.id.as_str()) {
errors.push(ValidationError::InvalidSchedule {
message: format!("duplicate schedule id {:?}", s.id),
});
}
match &s.trigger {
ScheduleTrigger::Manual => {}
ScheduleTrigger::Minutes { interval } | ScheduleTrigger::Hours { interval } => {
if *interval == 0 {
errors.push(ValidationError::InvalidSchedule {
message: format!("schedule {:?} interval must be positive", s.id),
});
}
}
ScheduleTrigger::Cron { cron } => {
if cron.trim().is_empty() {
errors.push(ValidationError::InvalidSchedule {
message: format!("schedule {:?} has an empty cron expression", s.id),
});
}
}
}
}
}
fn validate_requires(req: &crate::requires::Requires, errors: &mut Vec<ValidationError>) {
use crate::requires::{is_valid_pack_id, is_valid_sha256};
let mut seen: HashSet<&str> = HashSet::new();
for pr in &req.packs {
if !is_valid_pack_id(&pr.id) {
errors.push(ValidationError::InvalidRequires {
message: format!(
"pack id {:?} must match [a-z0-9][a-z0-9_-]{{0,63}}",
pr.id
),
});
}
if !seen.insert(pr.id.as_str()) {
errors.push(ValidationError::InvalidRequires {
message: format!("duplicate pack requirement {:?}", pr.id),
});
}
if let Some(range) = &pr.version
&& semver::VersionReq::parse(range).is_err()
{
errors.push(ValidationError::InvalidRequires {
message: format!("pack {:?} has unparseable version range {range:?}", pr.id),
});
}
if let Some(rv) = &pr.resolved_version
&& semver::Version::parse(rv).is_err()
{
errors.push(ValidationError::InvalidRequires {
message: format!("pack {:?} has unparseable resolved_version {rv:?}", pr.id),
});
}
if let Some(hash) = &pr.content_sha256
&& !is_valid_sha256(hash)
{
errors.push(ValidationError::InvalidRequires {
message: format!(
"pack {:?} content_sha256 must be 64 lowercase hex chars",
pr.id
),
});
}
}
}
pub fn validate_definition_only(def: &FlowDefinition) -> Vec<ValidationError> {
let mut errors = Vec::new();
validate_definition(def, crate::SPEC_VERSION, &mut errors);
errors
}
fn validate_definition(def: &FlowDefinition, spec_version: &str, errors: &mut Vec<ValidationError>) {
let entry_count = def
.nodes
.iter()
.filter(|n| matches!(n.node_type, FlowNodeType::Core(CoreNodeType::Entry)))
.count();
if entry_count > 1 {
errors.push(ValidationError::MultipleEntryNodes(entry_count));
}
let mut seen_nodes = HashSet::new();
for n in &def.nodes {
if !seen_nodes.insert(n.id.as_str()) {
errors.push(ValidationError::DuplicateNodeId(n.id.clone()));
}
match &n.node_type {
FlowNodeType::Custom(s) => {
if let Some((prefix, _)) = s.split_once(':') {
if !is_valid_vendor(prefix) {
errors.push(ValidationError::InvalidVendorNamespace {
node: n.id.clone(),
node_type: s.clone(),
});
}
} else {
errors.push(ValidationError::InvalidVendorNamespace {
node: n.id.clone(),
node_type: s.clone(),
});
}
}
FlowNodeType::Core(core) => {
if spec_version == "1" && core.is_v2() {
errors.push(ValidationError::V2NodeInV1Document {
node: n.id.clone(),
node_type: core.as_str().to_string(),
});
}
validate_core_node_data(n, *core, def, errors);
}
}
}
let node_ids: HashSet<&str> = def.nodes.iter().map(|n| n.id.as_str()).collect();
let mut seen_edges = HashSet::new();
for e in &def.edges {
if !seen_edges.insert(e.id.as_str()) {
errors.push(ValidationError::DuplicateEdgeId(e.id.clone()));
}
if !node_ids.contains(e.source.as_str()) {
errors.push(ValidationError::DanglingEdgeSource {
edge: e.id.clone(),
source: e.source.clone(),
});
}
if !node_ids.contains(e.target.as_str()) {
errors.push(ValidationError::DanglingEdgeTarget {
edge: e.id.clone(),
target: e.target.clone(),
});
}
}
}
fn has_handle_edge(def: &FlowDefinition, node_id: &str, handle: &str) -> bool {
def.edges
.iter()
.any(|e| e.source == node_id && e.source_handle.as_deref() == Some(handle))
}
fn validate_core_node_data(
node: &FlowNode,
core: CoreNodeType,
def: &FlowDefinition,
errors: &mut Vec<ValidationError>,
) {
match core {
CoreNodeType::Conditional => {
match serde_json::from_value::<ConditionalData>(node.data.clone()) {
Ok(data) => {
for cond in &data.conditions {
if Operator::from_wire(&cond.operator).is_none() {
errors.push(ValidationError::UnknownOperator {
node: node.id.clone(),
operator: cond.operator.clone(),
});
}
if !has_handle_edge(def, &node.id, &cond.handle) {
errors.push(ValidationError::MissingHandleEdge {
node: node.id.clone(),
handle: cond.handle.clone(),
});
}
}
if let Some(dh) = &data.default_handle
&& !has_handle_edge(def, &node.id, dh)
{
errors.push(ValidationError::MissingHandleEdge {
node: node.id.clone(),
handle: dh.clone(),
});
}
}
Err(e) => errors.push(ValidationError::InvalidNodeData {
node: node.id.clone(),
message: format!("expected conditional data: {e}"),
}),
}
}
CoreNodeType::Branch => {
match serde_json::from_value::<BranchData>(node.data.clone()) {
Ok(data) => {
if data.outputs.is_empty() {
errors.push(ValidationError::InvalidNodeData {
node: node.id.clone(),
message: "branch must declare at least one output".into(),
});
}
for out in &data.outputs {
if !has_handle_edge(def, &node.id, &out.handle) {
errors.push(ValidationError::MissingHandleEdge {
node: node.id.clone(),
handle: out.handle.clone(),
});
}
if out.handle == crate::BRANCH_ERROR_HANDLE
&& out.schema.as_ref().is_some_and(|s| {
s.get("type").and_then(|t| t.as_str()) != Some("string")
})
{
errors.push(ValidationError::InvalidNodeData {
node: node.id.clone(),
message: "the reserved `error` handle carries a string reason; \
its schema must be omitted or {\"type\":\"string\"}"
.into(),
});
}
}
if let Some(dh) = &data.default_handle
&& !has_handle_edge(def, &node.id, dh)
{
errors.push(ValidationError::MissingHandleEdge {
node: node.id.clone(),
handle: dh.clone(),
});
}
}
Err(e) => errors.push(ValidationError::InvalidNodeData {
node: node.id.clone(),
message: format!("expected branch data: {e}"),
}),
}
}
_ => {}
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::model::{FlowEdge, FlowNode};
use serde_json::json;
fn entry(id: &str) -> FlowNode {
FlowNode {
id: id.into(),
node_type: FlowNodeType::Core(CoreNodeType::Entry),
data: json!({}),
position: [0.0, 0.0],
}
}
fn prompt(id: &str) -> FlowNode {
FlowNode {
id: id.into(),
node_type: FlowNodeType::Core(CoreNodeType::Prompt),
data: json!({}),
position: [0.0, 0.0],
}
}
fn edge(id: &str, src: &str, tgt: &str) -> FlowEdge {
FlowEdge {
id: id.into(),
source: src.into(),
target: tgt.into(),
source_handle: None,
target_handle: None,
}
}
fn saved(def: FlowDefinition) -> SavedFlow {
SavedFlow {
spec_version: "1".into(),
id: "ok-id".into(),
name: "X".into(),
created_at: "2026-01-01T00:00:00Z".into(),
updated_at: "2026-01-01T00:00:00Z".into(),
enabled: false,
schedules: vec![],
requires: None,
flow: def,
}
}
#[test]
fn valid_minimal_flow_has_no_errors() {
let def = FlowDefinition {
nodes: vec![entry("e")],
edges: vec![],
};
assert!(validate(&saved(def)).is_empty());
}
#[test]
fn invalid_flow_id_caught() {
let mut sf = saved(FlowDefinition::default());
sf.id = "bad id with spaces".into();
let errs = validate(&sf);
assert!(errs.iter().any(|e| matches!(e, ValidationError::InvalidFlowId(_))));
}
#[test]
fn multiple_entries_caught() {
let def = FlowDefinition {
nodes: vec![entry("a"), entry("b")],
edges: vec![],
};
let errs = validate(&saved(def));
assert!(errs
.iter()
.any(|e| matches!(e, ValidationError::MultipleEntryNodes(2))));
}
#[test]
fn dangling_edge_caught() {
let def = FlowDefinition {
nodes: vec![entry("e")],
edges: vec![edge("x", "e", "missing")],
};
let errs = validate(&saved(def));
assert!(errs
.iter()
.any(|e| matches!(e, ValidationError::DanglingEdgeTarget { .. })));
}
#[test]
fn duplicate_node_id_caught() {
let def = FlowDefinition {
nodes: vec![entry("e"), prompt("e")],
edges: vec![],
};
let errs = validate(&saved(def));
assert!(errs.iter().any(|e| matches!(e, ValidationError::DuplicateNodeId(_))));
}
#[test]
fn both_spec_versions_accepted() {
let mut sf = saved(FlowDefinition {
nodes: vec![entry("e")],
edges: vec![],
});
sf.spec_version = "1".into();
assert!(validate(&sf).is_empty());
sf.spec_version = "2".into();
assert!(validate(&sf).is_empty());
}
#[test]
fn unsupported_spec_version_caught() {
let mut sf = saved(FlowDefinition::default());
sf.spec_version = "3".into();
let errs = validate(&sf);
assert!(errs
.iter()
.any(|e| matches!(e, ValidationError::UnsupportedSpecVersion(_))));
}
fn core(id: &str, ty: CoreNodeType, data: serde_json::Value) -> FlowNode {
FlowNode {
id: id.into(),
node_type: FlowNodeType::Core(ty),
data,
position: [0.0, 0.0],
}
}
fn eh(id: &str, src: &str, tgt: &str, handle: &str) -> FlowEdge {
FlowEdge {
id: id.into(),
source: src.into(),
target: tgt.into(),
source_handle: Some(handle.into()),
target_handle: None,
}
}
#[test]
fn v2_node_in_v1_document_caught() {
let def = FlowDefinition {
nodes: vec![entry("e"), core("c", CoreNodeType::Conditional, json!({ "conditions": [] }))],
edges: vec![edge("x", "e", "c")],
};
let mut sf = saved(def);
sf.spec_version = "1".into();
let errs = validate(&sf);
assert!(errs.iter().any(|e| matches!(e, ValidationError::V2NodeInV1Document { .. })));
}
#[test]
fn valid_conditional_v2_passes() {
let def = FlowDefinition {
nodes: vec![
entry("e"),
core("c", CoreNodeType::Conditional, json!({
"conditions": [ { "handle": "hot", "variable": "_last", "operator": "gt", "value": 50 } ],
"default_handle": "cold"
})),
prompt("hot_node"),
prompt("cold_node"),
],
edges: vec![
edge("e0", "e", "c"),
eh("e1", "c", "hot_node", "hot"),
eh("e2", "c", "cold_node", "cold"),
],
};
let mut sf = saved(def);
sf.spec_version = "2".into();
assert!(validate(&sf).is_empty(), "{:?}", validate(&sf));
}
#[test]
fn conditional_unknown_operator_and_missing_edge_caught() {
let def = FlowDefinition {
nodes: vec![
entry("e"),
core("c", CoreNodeType::Conditional, json!({
"conditions": [ { "handle": "hot", "variable": "_last", "operator": "bogus", "value": 1 } ]
})),
],
edges: vec![edge("e0", "e", "c")], };
let mut sf = saved(def);
sf.spec_version = "2".into();
let errs = validate(&sf);
assert!(errs.iter().any(|e| matches!(e, ValidationError::UnknownOperator { .. })));
assert!(errs.iter().any(|e| matches!(e, ValidationError::MissingHandleEdge { .. })));
}
#[test]
fn branch_bad_data_caught() {
let def = FlowDefinition {
nodes: vec![
entry("e"),
core("b", CoreNodeType::Branch, json!({ "persona": "weather-agent" })),
],
edges: vec![edge("e0", "e", "b")],
};
let mut sf = saved(def);
sf.spec_version = "2".into();
let errs = validate(&sf);
assert!(errs.iter().any(|e| matches!(e, ValidationError::InvalidNodeData { .. })));
}
#[test]
fn branch_error_handle_with_nonstring_schema_caught() {
let def = FlowDefinition {
nodes: vec![
entry("e"),
core("b", CoreNodeType::Branch, json!({
"query": "classify",
"outputs": [
{ "handle": "ok", "schema": { "type": "string" } },
{ "handle": "error", "schema": { "type": "object" } }
]
})),
prompt("ok_t"),
prompt("err_t"),
],
edges: vec![
edge("e0", "e", "b"),
eh("e1", "b", "ok_t", "ok"),
eh("e2", "b", "err_t", "error"),
],
};
let mut sf = saved(def);
sf.spec_version = "2".into();
let errs = validate(&sf);
assert!(
errs.iter().any(|e| matches!(
e,
ValidationError::InvalidNodeData { node, message }
if node == "b" && message.contains("`error` handle")
)),
"expected reserved-error-handle error, got {errs:?}"
);
}
#[test]
fn branch_error_handle_with_string_schema_passes() {
let def = FlowDefinition {
nodes: vec![
entry("e"),
core("b", CoreNodeType::Branch, json!({
"query": "classify",
"outputs": [
{ "handle": "ok", "schema": { "type": "string" } },
{ "handle": "error", "schema": { "type": "string" } }
]
})),
prompt("ok_t"),
prompt("err_t"),
],
edges: vec![
edge("e0", "e", "b"),
eh("e1", "b", "ok_t", "ok"),
eh("e2", "b", "err_t", "error"),
],
};
let mut sf = saved(def);
sf.spec_version = "2".into();
assert!(validate(&sf).is_empty(), "{:?}", validate(&sf));
}
#[test]
fn well_formed_custom_type_passes() {
let mut p = prompt("p");
p.node_type = FlowNodeType::Custom("slack:send_message".into());
let def = FlowDefinition {
nodes: vec![entry("e"), p],
edges: vec![edge("x", "e", "p")],
};
assert!(validate(&saved(def)).is_empty());
}
#[test]
fn malformed_custom_type_caught() {
let mut p = prompt("p");
p.node_type = FlowNodeType::Custom("BadVendor:thing".into());
let def = FlowDefinition {
nodes: vec![entry("e"), p],
edges: vec![],
};
let errs = validate(&saved(def));
assert!(errs
.iter()
.any(|e| matches!(e, ValidationError::InvalidVendorNamespace { .. })));
}
#[test]
fn custom_type_without_colon_caught() {
let mut p = prompt("p");
p.node_type = FlowNodeType::Custom("no_namespace".into());
let def = FlowDefinition {
nodes: vec![entry("e"), p],
edges: vec![],
};
let errs = validate(&saved(def));
assert!(errs
.iter()
.any(|e| matches!(e, ValidationError::InvalidVendorNamespace { .. })));
}
#[test]
fn well_formed_requires_passes() {
let mut sf = saved(FlowDefinition {
nodes: vec![entry("e")],
edges: vec![],
});
sf.requires = Some(crate::requires::Requires {
packs: vec![crate::requires::PackRequirement {
id: "cloudflare".into(),
version: Some(">=1.2.0, <2.0.0".into()),
content_sha256: Some("a1b2c3d4".repeat(8)),
resolved_version: Some("1.3.1".into()),
..crate::requires::PackRequirement::new("cloudflare")
}],
tools: vec!["cloudflare_purge_cache".into()],
});
assert!(validate(&sf).is_empty(), "{:?}", validate(&sf));
}
#[test]
fn malformed_requires_caught() {
for req in [
crate::requires::Requires {
packs: vec![crate::requires::PackRequirement::new("Bad Id")],
tools: vec![],
},
crate::requires::Requires {
packs: vec![crate::requires::PackRequirement {
version: Some("not-a-range".into()),
..crate::requires::PackRequirement::new("cloudflare")
}],
tools: vec![],
},
crate::requires::Requires {
packs: vec![crate::requires::PackRequirement {
content_sha256: Some("tooshort".into()),
..crate::requires::PackRequirement::new("cloudflare")
}],
tools: vec![],
},
crate::requires::Requires {
packs: vec![
crate::requires::PackRequirement::new("cloudflare"),
crate::requires::PackRequirement::new("cloudflare"),
],
tools: vec![],
},
] {
let mut sf = saved(FlowDefinition {
nodes: vec![entry("e")],
edges: vec![],
});
sf.requires = Some(req);
let errs = validate(&sf);
assert!(
errs.iter()
.any(|e| matches!(e, ValidationError::InvalidRequires { .. })),
"expected InvalidRequires, got {errs:?}"
);
}
}
}