use std::collections::BTreeMap;
use serde::{Deserialize, Serialize};
use std::path::Path;
use crate::engine::execution::{node_key, ExecutionCursor, NestedCursor};
use crate::engine::flow::ConcreteStep;
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct FlowGraph {
pub name: String,
pub steps: Vec<FlowNode>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct FlowNode {
pub key: String,
pub id: Option<String>,
pub label: String,
pub kind: FlowNodeKind,
pub human: bool,
pub returns_to: Option<String>,
pub parents: Vec<String>,
pub paths: Vec<FlowGraphPath>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum FlowNodeKind {
Skill,
Op,
Xor,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct FlowGraphPath {
pub name: String,
pub description: String,
pub steps: Vec<FlowNode>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct FlowReturn {
pub decider: String,
pub traversals: u32,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct FlowCatalogEntry {
pub name: String,
pub graph: Option<FlowGraph>,
pub unavailable: Option<String>,
}
pub fn flow_catalog(repo: &Path) -> Vec<FlowCatalogEntry> {
crate::engine::available_flow_names(repo)
.into_iter()
.map(|name| {
let expanded = crate::engine::load_flow(&name, repo)
.map_err(|error| error.to_string())
.and_then(|flow| {
crate::engine::expand_flow(&flow, repo)
.map(|steps| FlowGraph::new(&flow.name, &steps))
.map_err(|error| error.to_string())
});
match expanded {
Ok(graph) => FlowCatalogEntry {
name,
graph: Some(graph),
unavailable: None,
},
Err(reason) => FlowCatalogEntry {
name,
graph: None,
unavailable: Some(reason),
},
}
})
.collect()
}
impl FlowGraph {
pub fn new(name: impl Into<String>, steps: &[ConcreteStep]) -> Self {
Self {
name: name.into(),
steps: nodes(steps, ""),
}
}
}
fn nodes(steps: &[ConcreteStep], prefix: &str) -> Vec<FlowNode> {
steps
.iter()
.enumerate()
.map(|(index, step)| {
let key = node_key(prefix, index);
match step {
ConcreteStep::Skill(skill) => FlowNode {
returns_to: return_target(steps, index).map(|target| node_key(prefix, target)),
key,
id: skill.id.clone(),
label: skill.skill.name.clone(),
kind: FlowNodeKind::Skill,
human: skill.human,
parents: skill.flow_parents.clone(),
paths: Vec::new(),
},
ConcreteStep::Command(op) => FlowNode {
key,
id: None,
label: op.item.display_name(),
kind: FlowNodeKind::Op,
human: false,
returns_to: None,
parents: op.flow_parents.clone(),
paths: Vec::new(),
},
ConcreteStep::Xor(branch) => {
let mut names: Vec<_> = branch.paths.keys().collect();
names.sort();
let paths = names
.into_iter()
.map(|name| {
let path = &branch.paths[name];
FlowGraphPath {
name: name.clone(),
description: path.description.clone(),
steps: nodes(&path.steps, &format!("{key}/{name}/")),
}
})
.collect();
FlowNode {
key,
id: None,
label: branch.router.name.clone(),
kind: FlowNodeKind::Xor,
human: false,
returns_to: None,
parents: branch.flow_parents.clone(),
paths,
}
}
}
})
.collect()
}
fn return_target(steps: &[ConcreteStep], index: usize) -> Option<usize> {
let ConcreteStep::Skill(skill) = &steps[index] else {
return None;
};
let from = skill.repeat.as_ref()?.from.as_str();
steps[..index].iter().position(|target| {
matches!(target, ConcreteStep::Skill(target) if target.id.as_deref() == Some(from))
})
}
pub fn flow_iterations(steps: &[ConcreteStep], cursor: &ExecutionCursor) -> Vec<Vec<u32>> {
let counts = steps
.iter()
.enumerate()
.filter_map(|(index, step)| {
return_target(steps, index)?;
let ConcreteStep::Skill(skill) = step else {
return None;
};
let id = skill.id.as_ref()?;
Some(cursor.progress.repeats.get(id).copied().unwrap_or(0))
})
.collect();
let mut levels = vec![counts];
if let (
Some(ConcreteStep::Xor(branch)),
Some(NestedCursor::Xor {
selected,
cursor: child,
}),
) = (steps.get(cursor.index), cursor.child.as_deref())
{
if let Some(path) = branch.paths.get(selected) {
levels.extend(flow_iterations(&path.steps, child));
}
}
levels
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct CursorProjection {
pub current: Option<String>,
pub completed: Vec<String>,
pub returns: Vec<FlowReturn>,
}
pub fn project_cursor(steps: &[ConcreteStep], cursor: &ExecutionCursor) -> CursorProjection {
let mut projection = CursorProjection {
current: None,
completed: Vec::new(),
returns: Vec::new(),
};
walk_cursor(steps, "", Some(cursor), &mut projection);
collect_returns(
steps,
"",
Some(cursor),
&cursor.progress.repeats,
&mut projection.returns,
);
projection
}
fn walk_cursor(
steps: &[ConcreteStep],
prefix: &str,
cursor: Option<&ExecutionCursor>,
projection: &mut CursorProjection,
) {
let Some(cursor) = cursor else { return };
let index = cursor.index.min(steps.len());
projection
.completed
.extend((0..index).map(|index| node_key(prefix, index)));
let key = node_key(prefix, cursor.index);
match (steps.get(cursor.index), cursor.child.as_deref()) {
(Some(ConcreteStep::Xor(branch)), Some(NestedCursor::Xor { selected, cursor })) => {
match branch.paths.get(selected) {
Some(path) if cursor.index < path.steps.len() => {
projection.completed.push(key.clone());
walk_cursor(
&path.steps,
&format!("{key}/{selected}/"),
Some(cursor),
projection,
);
}
_ => projection.current = Some(key),
}
}
(Some(_), _) => projection.current = Some(key),
(None, _) => {}
}
}
fn collect_returns(
steps: &[ConcreteStep],
prefix: &str,
cursor: Option<&ExecutionCursor>,
repeats: &BTreeMap<String, u32>,
out: &mut Vec<FlowReturn>,
) {
for (index, step) in steps.iter().enumerate() {
let key = node_key(prefix, index);
match step {
ConcreteStep::Skill(skill) => {
let (Some(id), Some(_target)) = (&skill.id, return_target(steps, index)) else {
continue;
};
out.push(FlowReturn {
decider: key,
traversals: repeats.get(id).copied().unwrap_or(0),
});
}
ConcreteStep::Xor(branch) => {
let mut names: Vec<_> = branch.paths.keys().collect();
names.sort();
for name in names {
let child = cursor
.filter(|cursor| cursor.index == index)
.and_then(|cursor| match cursor.child.as_deref() {
Some(NestedCursor::Xor { selected, cursor }) if selected == name => {
Some(cursor)
}
_ => None,
});
let settled_prefix = format!("xor:{index}:{name}/");
let settled: BTreeMap<String, u32> = repeats
.iter()
.filter_map(|(key, count)| {
key.strip_prefix(&settled_prefix)
.map(|key| (key.to_owned(), *count))
})
.collect();
collect_returns(
&branch.paths[name].steps,
&format!("{key}/{name}/"),
child,
child.map_or(&settled, |child| &child.progress.repeats),
out,
);
}
}
ConcreteStep::Command(_) => {}
}
}
}
#[cfg(test)]
mod tests {
use std::collections::{BTreeMap, HashMap};
use crate::engine::execution::{ExecutionCursor, NestedCursor};
use crate::engine::flow::{
ConcretePath, ConcreteSkill, ConcreteStep, ConcreteXor, RepeatPolicy, Skill,
};
use crate::engine::flow_graph::{flow_iterations, project_cursor, FlowGraph, FlowNodeKind};
use crate::engine::{expand_flow, load_flow};
fn skill(name: &str, id: Option<&str>, human: bool, from: Option<&str>) -> ConcreteStep {
ConcreteStep::Skill(ConcreteSkill {
skill: Skill::named(name),
id: id.map(str::to_string),
human,
repeat: from.map(|from| RepeatPolicy {
from: from.to_string(),
}),
flow_parents: Vec::new(),
})
}
#[test]
fn retained_flows_keep_their_delivery_and_review_boundaries() {
let repo = tempfile::tempdir().unwrap();
let cases: &[(&str, &[&str], &[usize], usize)] = &[
("code", &["implement", "compress"], &[], 0),
("queue", &["compress", "sync", "realign", "gate"], &[], 0),
("refresh", &["sync", "realign"], &[], 0),
("task-design", &["kickoff", "review-design"], &[1], 0),
("incident", &["unbreak", "5whys", "launch-plan"], &[], 0),
(
"pursue",
&[
"implement",
"compress",
"sync",
"realign",
"loop-decide",
"pr-publish",
"demo",
"loop-decide",
],
&[6],
2,
),
("deploy", &["gate", "pr land"], &[], 0),
("ship", &["gate", "pr land -c"], &[], 0),
("ship-demo", &["gate", "demo", "pr land -c"], &[1], 0),
("vsm-operate", &["s1", "s2", "s3", "s4", "s5"], &[], 0),
];
for (name, labels, humans, returns) in cases {
let flow = load_flow(name, repo.path()).unwrap();
let graph = FlowGraph::new(*name, &expand_flow(&flow, repo.path()).unwrap());
assert_eq!(
graph
.steps
.iter()
.map(|s| s.label.as_str())
.collect::<Vec<_>>(),
*labels,
"{name}"
);
assert_eq!(
graph
.steps
.iter()
.enumerate()
.filter_map(|(i, s)| s.human.then_some(i))
.collect::<Vec<_>>(),
*humans,
"{name}"
);
assert_eq!(
graph
.steps
.iter()
.filter(|s| s.returns_to.is_some())
.count(),
*returns,
"{name}"
);
}
}
#[test]
fn refresh_integrates_upstream_before_realigning() {
let repo = tempfile::tempdir().unwrap();
let flow = load_flow("refresh", repo.path()).unwrap();
let graph = FlowGraph::new(&flow.name, &expand_flow(&flow, repo.path()).unwrap());
let labels: Vec<_> = graph.steps.iter().map(|node| node.label.as_str()).collect();
assert_eq!(labels, ["sync", "realign"]);
assert_eq!(graph.steps[0].kind, FlowNodeKind::Op);
assert_eq!(graph.steps[1].kind, FlowNodeKind::Skill);
}
#[test]
fn feature_draws_both_returns_to_implement_with_forward_delivery() {
let repo = tempfile::tempdir().unwrap();
let flow = load_flow("feature", repo.path()).unwrap();
let graph = FlowGraph::new(&flow.name, &expand_flow(&flow, repo.path()).unwrap());
let labels: Vec<_> = graph.steps.iter().map(|node| node.label.as_str()).collect();
assert_eq!(
labels,
[
"kickoff",
"review-design",
"implement",
"compress",
"sync",
"realign",
"loop-decide",
"pr-publish",
"demo",
"loop-decide",
"compress",
"sync",
"realign",
"gate",
"pr land -c"
]
);
let implement = graph.steps[2].key.clone();
let returns: Vec<_> = graph
.steps
.iter()
.filter_map(|node| node.returns_to.as_ref().map(|to| (node.id.clone(), to)))
.collect();
assert_eq!(
returns,
[
(Some("decide".to_string()), &implement),
(Some("decide_delivery".to_string()), &implement)
]
);
assert!(graph.steps[1].human && graph.steps[8].human);
assert_eq!(graph.steps[14].kind, FlowNodeKind::Op);
assert_eq!(graph.steps[5].parents, ["feature", "pursue", "refresh"]);
}
#[test]
fn repeated_decisions_keep_independent_counts_and_a_pass_scoped_completion() {
let steps = vec![
skill("design", Some("design"), false, None),
skill("implement", Some("implement"), false, None),
skill("loop-decide", Some("decide"), false, Some("implement")),
skill("demo", Some("demo"), true, None),
skill(
"loop-decide",
Some("decide_delivery"),
false,
Some("implement"),
),
skill("land", None, false, None),
];
let cursor = ExecutionCursor {
index: 1,
iteration: 3,
progress: crate::engine::transitions::FlowProgress {
repeats: BTreeMap::from([("decide".into(), 2), ("decide_delivery".into(), 1)]),
..Default::default()
},
..Default::default()
};
let projection = project_cursor(&steps, &cursor);
assert_eq!(projection.current.as_deref(), Some("1"));
assert_eq!(projection.completed, ["0"]);
let counts: Vec<_> = projection
.returns
.iter()
.map(|edge| (edge.decider.as_str(), edge.traversals))
.collect();
assert_eq!(counts, [("2", 2), ("4", 1)]);
assert_eq!(flow_iterations(&steps, &cursor), [vec![2, 1]]);
let later = ExecutionCursor {
index: 3,
iteration: 99,
..cursor
};
assert_eq!(flow_iterations(&steps, &later), [vec![2, 1]]);
}
#[test]
fn a_selected_xor_path_is_drawn_honestly_with_nested_keys() {
let branch = ConcreteXor {
router: Skill::named("xor-route"),
paths: HashMap::from([
(
"fix".to_string(),
ConcretePath {
description: "Fix it".into(),
steps: vec![
skill("patch", Some("patch"), false, None),
skill("check", Some("check"), false, Some("patch")),
],
},
),
(
"skip".to_string(),
ConcretePath {
description: "Nothing to do".into(),
steps: Vec::new(),
},
),
]),
flow_parents: Vec::new(),
};
let steps = vec![
skill("kickoff", None, false, None),
ConcreteStep::Xor(branch),
];
let graph = FlowGraph::new("routed", &steps);
let xor = &graph.steps[1];
assert_eq!(xor.kind, FlowNodeKind::Xor);
assert_eq!(
xor.paths
.iter()
.map(|path| path.name.as_str())
.collect::<Vec<_>>(),
["fix", "skip"]
);
assert_eq!(xor.paths[0].steps[1].key, "1/fix/1");
assert_eq!(xor.paths[0].steps[1].returns_to.as_deref(), Some("1/fix/0"));
let cursor = ExecutionCursor {
index: 1,
progress: crate::engine::transitions::FlowProgress {
repeats: BTreeMap::from([("xor:1:fix/check".into(), 4)]),
..Default::default()
},
child: Some(Box::new(NestedCursor::Xor {
selected: "fix".into(),
cursor: ExecutionCursor {
index: 1,
iteration: 2,
progress: crate::engine::transitions::FlowProgress {
repeats: BTreeMap::from([("check".into(), 2)]),
..Default::default()
},
..Default::default()
},
})),
..Default::default()
};
let projection = project_cursor(&steps, &cursor);
assert_eq!(projection.current.as_deref(), Some("1/fix/1"));
assert_eq!(projection.completed, ["0", "1", "1/fix/0"]);
assert_eq!(projection.returns[0].traversals, 2);
assert_eq!(flow_iterations(&steps, &cursor), [vec![], vec![2]]);
}
}