use serde::{Deserialize, Deserializer, Serialize, Serializer};
pub const SPEC_VERSION: &str = "2";
pub const SUPPORTED_SPEC_VERSIONS: &[&str] = &["1", "2"];
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
pub struct FlowNode {
pub id: String,
pub node_type: FlowNodeType,
pub data: serde_json::Value,
#[serde(default)]
pub position: [f64; 2],
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum FlowNodeType {
Core(CoreNodeType),
Custom(String),
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
pub enum CoreNodeType {
Entry,
Prompt,
Conditional,
Branch,
SetVariable,
Tool,
Http,
SubAgent,
Approval,
Wait,
Foreach,
End,
BranchTool,
}
impl CoreNodeType {
pub fn as_str(self) -> &'static str {
match self {
CoreNodeType::Entry => "entry",
CoreNodeType::Prompt => "prompt",
CoreNodeType::Conditional => "conditional",
CoreNodeType::Branch => "branch",
CoreNodeType::SetVariable => "set_variable",
CoreNodeType::Tool => "tool",
CoreNodeType::Http => "http",
CoreNodeType::SubAgent => "sub_agent",
CoreNodeType::Approval => "approval",
CoreNodeType::Wait => "wait",
CoreNodeType::Foreach => "foreach",
CoreNodeType::End => "end",
CoreNodeType::BranchTool => "branch_tool",
}
}
pub fn from_wire(s: &str) -> Option<Self> {
match s {
"entry" => Some(CoreNodeType::Entry),
"prompt" => Some(CoreNodeType::Prompt),
"conditional" => Some(CoreNodeType::Conditional),
"branch" => Some(CoreNodeType::Branch),
"set_variable" => Some(CoreNodeType::SetVariable),
"tool" => Some(CoreNodeType::Tool),
"http" => Some(CoreNodeType::Http),
"sub_agent" => Some(CoreNodeType::SubAgent),
"approval" => Some(CoreNodeType::Approval),
"wait" => Some(CoreNodeType::Wait),
"foreach" => Some(CoreNodeType::Foreach),
"end" => Some(CoreNodeType::End),
"branch_tool" => Some(CoreNodeType::BranchTool),
_ => None,
}
}
pub fn is_v2(self) -> bool {
!matches!(
self,
CoreNodeType::Entry
| CoreNodeType::Prompt
| CoreNodeType::BranchTool
)
}
}
impl FlowNodeType {
pub fn as_wire(&self) -> &str {
match self {
FlowNodeType::Core(c) => c.as_str(),
FlowNodeType::Custom(s) => s.as_str(),
}
}
}
impl Serialize for FlowNodeType {
fn serialize<S: Serializer>(&self, ser: S) -> Result<S::Ok, S::Error> {
ser.serialize_str(self.as_wire())
}
}
impl<'de> Deserialize<'de> for FlowNodeType {
fn deserialize<D: Deserializer<'de>>(de: D) -> Result<Self, D::Error> {
let s = String::deserialize(de)?;
if let Some(core) = CoreNodeType::from_wire(&s) {
Ok(FlowNodeType::Core(core))
} else {
Ok(FlowNodeType::Custom(s))
}
}
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
pub struct FlowEdge {
pub id: String,
pub source: String,
pub target: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub source_handle: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub target_handle: Option<String>,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Default)]
pub struct FlowDefinition {
pub nodes: Vec<FlowNode>,
pub edges: Vec<FlowEdge>,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
pub struct SavedFlow {
#[serde(default = "default_spec_version")]
pub spec_version: String,
pub id: String,
pub name: String,
pub created_at: String,
pub updated_at: String,
#[serde(default)]
pub enabled: bool,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub schedules: Vec<FlowScheduleSpec>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub requires: Option<crate::requires::Requires>,
pub flow: FlowDefinition,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
pub struct FlowScheduleSpec {
pub id: String,
#[serde(default = "default_true")]
pub enabled: bool,
#[serde(flatten)]
pub trigger: ScheduleTrigger,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub name: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub timezone: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub inputs: Option<serde_json::Value>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub persona: Option<String>,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
#[serde(tag = "type", rename_all = "snake_case")]
pub enum ScheduleTrigger {
Manual,
Minutes {
interval: u64,
},
Hours {
interval: u64,
},
Cron {
cron: String,
},
}
fn default_true() -> bool {
true
}
impl SavedFlow {
pub fn effective_schedules(&self) -> Vec<FlowScheduleSpec> {
if !self.schedules.is_empty() {
return self.schedules.clone();
}
if let Some(spec) = self.entry_schedule_from_node() {
return vec![spec];
}
vec![FlowScheduleSpec {
id: "default".to_string(),
enabled: true,
trigger: ScheduleTrigger::Manual,
name: None,
timezone: None,
inputs: None,
persona: None,
}]
}
fn entry_schedule_from_node(&self) -> Option<FlowScheduleSpec> {
let entry = self
.flow
.nodes
.iter()
.find(|n| matches!(n.node_type, FlowNodeType::Core(CoreNodeType::Entry)))?;
let schedule_type = entry.data.get("schedule_type").and_then(|v| v.as_str())?;
let interval = entry
.data
.get("interval")
.and_then(|v| v.as_u64())
.unwrap_or(0);
let trigger = match schedule_type {
"minutes" => ScheduleTrigger::Minutes { interval },
"hours" => ScheduleTrigger::Hours { interval },
"cron" => ScheduleTrigger::Cron {
cron: entry
.data
.get("cron")
.and_then(|v| v.as_str())
.unwrap_or_default()
.to_string(),
},
_ => ScheduleTrigger::Manual,
};
let persona = entry
.data
.get("persona")
.and_then(|v| v.as_str())
.map(|s| s.to_string());
Some(FlowScheduleSpec {
id: "default".to_string(),
enabled: true,
trigger,
name: None,
timezone: None,
inputs: None,
persona,
})
}
}
pub const DEFAULT_SPEC_VERSION: &str = "1";
fn default_spec_version() -> String {
DEFAULT_SPEC_VERSION.to_string()
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
pub struct FlowSummary {
pub id: String,
pub name: String,
pub node_count: usize,
pub created_at: String,
pub updated_at: String,
#[serde(default)]
pub enabled: bool,
#[serde(default)]
pub schedule_count: usize,
}
pub(crate) fn is_safe_id(id: &str) -> bool {
!id.is_empty()
&& id.len() <= 64
&& id.chars().all(|c| c.is_ascii_alphanumeric() || c == '-')
}
pub(crate) fn is_valid_vendor(prefix: &str) -> bool {
let mut chars = prefix.chars();
let Some(first) = chars.next() else { return false };
if !first.is_ascii_lowercase() {
return false;
}
if prefix.len() > 32 {
return false;
}
chars.all(|c| c.is_ascii_lowercase() || c.is_ascii_digit() || c == '_' || c == '-')
}
#[cfg(test)]
mod tests {
use super::*;
use serde_json::json;
#[test]
fn core_node_type_round_trips() {
for ct in [
CoreNodeType::Entry,
CoreNodeType::Prompt,
CoreNodeType::Branch,
CoreNodeType::BranchTool,
] {
let nt = FlowNodeType::Core(ct);
let j = serde_json::to_string(&nt).unwrap();
let back: FlowNodeType = serde_json::from_str(&j).unwrap();
assert_eq!(nt, back);
}
}
#[test]
fn custom_node_type_round_trips() {
let nt = FlowNodeType::Custom("slack:send_message".to_string());
let j = serde_json::to_string(&nt).unwrap();
assert_eq!(j, "\"slack:send_message\"");
let back: FlowNodeType = serde_json::from_str(&j).unwrap();
assert_eq!(nt, back);
}
#[test]
fn unknown_bare_node_type_becomes_custom() {
let back: FlowNodeType = serde_json::from_str("\"future_core_type\"").unwrap();
assert_eq!(back, FlowNodeType::Custom("future_core_type".into()));
}
#[test]
fn missing_spec_version_defaults_to_v1() {
let doc = json!({
"id": "x",
"name": "X",
"created_at": "2026-01-01T00:00:00Z",
"updated_at": "2026-01-01T00:00:00Z",
"flow": { "nodes": [], "edges": [] }
});
let parsed: SavedFlow = serde_json::from_value(doc).unwrap();
assert_eq!(parsed.spec_version, "1");
assert!(!parsed.enabled);
}
#[test]
fn saved_flow_round_trips() {
let sf = SavedFlow {
spec_version: "1".into(),
id: "f1".into(),
name: "F1".into(),
created_at: "2026-01-01T00:00:00Z".into(),
updated_at: "2026-01-02T00:00:00Z".into(),
enabled: true,
schedules: vec![],
requires: None,
flow: FlowDefinition {
nodes: vec![FlowNode {
id: "n1".into(),
node_type: FlowNodeType::Core(CoreNodeType::Entry),
data: json!({"schedule_type": "manual"}),
position: [10.0, 20.0],
}],
edges: vec![],
},
};
let j = serde_json::to_string(&sf).unwrap();
let back: SavedFlow = serde_json::from_str(&j).unwrap();
assert_eq!(sf, back);
}
#[test]
fn effective_schedules_prefers_top_level_array() {
let mut sf: SavedFlow = serde_json::from_value(json!({
"id": "f", "name": "F",
"created_at": "2026-01-01T00:00:00Z", "updated_at": "2026-01-01T00:00:00Z",
"schedules": [
{ "id": "morning", "type": "cron", "cron": "0 8 * * *" },
{ "id": "evening", "type": "cron", "cron": "0 18 * * *", "enabled": false }
],
"flow": { "nodes": [
{ "id": "entry", "node_type": "entry", "data": { "schedule_type": "cron", "cron": "0 0 * * *" }, "position": [0,0] }
], "edges": [] }
}))
.unwrap();
let eff = sf.effective_schedules();
assert_eq!(eff.len(), 2, "top-level array wins over the entry node");
assert_eq!(eff[0].id, "morning");
assert!(eff[0].enabled);
assert!(!eff[1].enabled);
assert_eq!(eff[0].trigger, ScheduleTrigger::Cron { cron: "0 8 * * *".into() });
sf.schedules.clear();
let eff = sf.effective_schedules();
assert_eq!(eff.len(), 1);
assert_eq!(eff[0].trigger, ScheduleTrigger::Cron { cron: "0 0 * * *".into() });
}
#[test]
fn effective_schedules_legacy_entry_and_manual_fallback() {
let sf: SavedFlow = serde_json::from_value(json!({
"id": "f", "name": "F",
"created_at": "2026-01-01T00:00:00Z", "updated_at": "2026-01-01T00:00:00Z",
"flow": { "nodes": [
{ "id": "entry", "node_type": "entry", "data": {}, "position": [0,0] }
], "edges": [] }
}))
.unwrap();
let eff = sf.effective_schedules();
assert_eq!(eff.len(), 1);
assert_eq!(eff[0].trigger, ScheduleTrigger::Manual);
let sf: SavedFlow = serde_json::from_value(json!({
"id": "f", "name": "F",
"created_at": "2026-01-01T00:00:00Z", "updated_at": "2026-01-01T00:00:00Z",
"flow": { "nodes": [
{ "id": "entry", "node_type": "entry", "data": { "schedule_type": "minutes", "interval": 15, "persona": "briefer" }, "position": [0,0] }
], "edges": [] }
}))
.unwrap();
let eff = sf.effective_schedules();
assert_eq!(eff[0].trigger, ScheduleTrigger::Minutes { interval: 15 });
assert_eq!(eff[0].persona.as_deref(), Some("briefer"));
}
#[test]
fn schedule_trigger_serializes_with_type_tag() {
let spec = FlowScheduleSpec {
id: "s".into(),
enabled: true,
trigger: ScheduleTrigger::Cron { cron: "0 8 * * *".into() },
name: Some("Morning".into()),
timezone: Some("America/Detroit".into()),
inputs: None,
persona: None,
};
let v = serde_json::to_value(&spec).unwrap();
assert_eq!(v["type"], "cron");
assert_eq!(v["cron"], "0 8 * * *");
assert_eq!(v["timezone"], "America/Detroit");
let back: FlowScheduleSpec =
serde_json::from_value(json!({ "id": "s", "type": "manual" })).unwrap();
assert!(back.enabled);
}
#[test]
fn id_validation() {
assert!(is_safe_id("ok-id"));
assert!(is_safe_id("a"));
assert!(!is_safe_id(""));
assert!(!is_safe_id("has space"));
assert!(!is_safe_id("../escape"));
assert!(!is_safe_id(&"x".repeat(65)));
}
#[test]
fn vendor_validation() {
assert!(is_valid_vendor("slack"));
assert!(is_valid_vendor("my-co"));
assert!(is_valid_vendor("my_co"));
assert!(is_valid_vendor("co0"));
assert!(!is_valid_vendor(""));
assert!(!is_valid_vendor("0starts-with-digit"));
assert!(!is_valid_vendor("Capital"));
assert!(!is_valid_vendor(&"a".repeat(33)));
}
}