use std::collections::{BTreeMap, BTreeSet};
use serde::{Deserialize, Serialize};
use serde_json::Value;
use crate::core::{Capability, Digest, StepId, canon};
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum Topology {
#[default]
Single,
Collaborative(Collaboration),
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "kebab-case")]
pub enum Collaboration {
ParallelDisjoint,
DistinctAuthority,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
#[serde(tag = "from", rename_all = "snake_case")]
pub enum ArgSource {
RunInput {
#[serde(default, skip_serializing_if = "Option::is_none")]
field: Option<String>,
},
Node {
step: StepId,
#[serde(default, skip_serializing_if = "Option::is_none")]
field: Option<String>,
},
Const { value: Value },
}
impl ArgSource {
#[must_use]
pub fn run_input() -> Self {
Self::RunInput { field: None }
}
#[must_use]
pub fn input_field(f: impl Into<String>) -> Self {
Self::RunInput {
field: Some(f.into()),
}
}
#[must_use]
pub fn node(step: StepId) -> Self {
Self::Node { step, field: None }
}
#[must_use]
pub fn node_field(step: StepId, f: impl Into<String>) -> Self {
Self::Node {
step,
field: Some(f.into()),
}
}
#[must_use]
pub fn constant(v: Value) -> Self {
Self::Const { value: v }
}
#[must_use]
pub fn depends_on(&self) -> Option<StepId> {
match self {
Self::Node { step, .. } => Some(*step),
_ => None,
}
}
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct PlanNode {
pub id: StepId,
pub capability: Capability,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub depends_on: Vec<StepId>,
#[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
pub args: BTreeMap<String, ArgSource>,
#[serde(default)]
pub terminal: bool,
#[serde(default)]
pub verifies: bool,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub quorum: Option<crate::core::Quorum>,
}
impl PlanNode {
pub fn new(id: u32, capability: impl Into<Capability>) -> Self {
Self {
quorum: None,
id: StepId(id),
capability: capability.into(),
depends_on: Vec::new(),
args: BTreeMap::new(),
terminal: false,
verifies: false,
}
}
#[must_use]
pub fn arg(mut self, name: impl Into<String>, source: ArgSource) -> Self {
if let Some(dep) = source.depends_on()
&& !self.depends_on.contains(&dep)
{
self.depends_on.push(dep);
}
self.args.insert(name.into(), source);
self
}
#[must_use]
pub fn after(mut self, step: u32) -> Self {
let s = StepId(step);
if !self.depends_on.contains(&s) {
self.depends_on.push(s);
}
self
}
#[must_use]
pub fn terminal(mut self) -> Self {
self.terminal = true;
self
}
#[must_use]
pub fn verifies(mut self) -> Self {
self.verifies = true;
self
}
#[must_use]
pub fn with_quorum(mut self, quorum: crate::core::Quorum) -> Self {
self.quorum = Some(quorum);
self
}
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct PlanIR {
pub version: u32,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub derived_from: Option<Digest>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub reason: Option<String>,
pub topology: Topology,
pub nodes: Vec<PlanNode>,
}
impl PlanIR {
#[must_use]
pub fn new(nodes: Vec<PlanNode>) -> Self {
Self {
version: 1,
derived_from: None,
reason: None,
topology: Topology::Single,
nodes,
}
}
#[must_use]
pub fn single(capability: impl Into<Capability>) -> Self {
Self::new(vec![
PlanNode::new(0, capability)
.arg("input", ArgSource::run_input())
.terminal(),
])
}
#[must_use]
pub fn fan_out(
branches: impl IntoIterator<Item = impl Into<Capability>>,
aggregate: impl Into<Capability>,
) -> Self {
let branches: Vec<Capability> = branches.into_iter().map(Into::into).collect();
assert!(
!branches.is_empty(),
"a fan-out needs at least one branch; with none the aggregator has \
nothing to aggregate"
);
let mut nodes: Vec<PlanNode> = branches
.iter()
.enumerate()
.map(|(i, capability)| {
PlanNode::new(u32::try_from(i).unwrap_or(u32::MAX), capability.clone())
.arg("input", ArgSource::run_input())
})
.collect();
let join_id = u32::try_from(branches.len()).unwrap_or(u32::MAX);
let mut join = PlanNode::new(join_id, aggregate).terminal();
for (i, capability) in branches.iter().enumerate() {
let step = StepId(u32::try_from(i).unwrap_or(u32::MAX));
join = join.arg(&capability.0, ArgSource::node(step));
}
nodes.push(join);
Self::new(nodes)
}
#[must_use]
pub fn topology(mut self, t: Topology) -> Self {
self.topology = t;
self
}
#[must_use]
pub fn digest(&self) -> Digest {
Digest::of(&canon::to_bytes(self).unwrap_or_default())
}
#[must_use]
pub fn node(&self, id: StepId) -> Option<&PlanNode> {
self.nodes.iter().find(|n| n.id == id)
}
#[must_use]
pub fn ready(&self, done: &BTreeSet<StepId>) -> Vec<StepId> {
let mut ready: Vec<StepId> = self
.nodes
.iter()
.filter(|n| !done.contains(&n.id))
.filter(|n| n.depends_on.iter().all(|d| done.contains(d)))
.map(|n| n.id)
.collect();
ready.sort_unstable();
ready
}
#[must_use]
pub fn is_complete(&self, done: &BTreeSet<StepId>) -> bool {
self.nodes
.iter()
.filter(|n| n.terminal)
.all(|n| done.contains(&n.id))
}
}
#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
#[non_exhaustive]
pub enum PlanError {
#[error("the plan has no nodes")]
Empty,
#[error("step {0} appears more than once")]
DuplicateStep(StepId),
#[error("step {step} depends on {missing}, which is not in the plan")]
MissingDependency { step: StepId, missing: StepId },
#[error("the plan has a dependency cycle involving {0}")]
Cycle(StepId),
#[error("the plan has no terminal node, so nothing marks it complete")]
NoTerminal,
#[error("step {step} is unreachable: nothing depends on it and it is not terminal")]
Unreachable { step: StepId },
#[error("no skill provides capability '{capability}' required by step {step}")]
NoProvider { step: StepId, capability: String },
#[error("step {step} takes argument '{arg}' from {from_step}, which is not an upstream node")]
ArgumentNotUpstream {
step: StepId,
arg: String,
from_step: StepId,
},
#[error("step {step} has no bound arguments, so its input is undefined")]
NoArguments { step: StepId },
#[error("step {step} verifies nothing: a verifier must depend on what it checks")]
VerifierWithoutSubject { step: StepId },
#[error(
"step {step} declares a quorum but judges nothing; a panel needs \
something to judge, or it is repetition rather than review"
)]
QuorumWithoutSubject { step: StepId },
#[error("this plan requires a verifier node and has none")]
VerifierRequired,
#[error(
"collaboration claims parallel-disjoint, but steps {a} and {b} read the same source — \
paying coordination cost for parallelism that is not there"
)]
FalseParallelism { a: StepId, b: StepId },
#[error(
"collaboration claims distinct-authority, but every step needs the same capability \
'{capability}' — there is no authority to separate"
)]
NoAuthorityToSeparate { capability: String },
#[error("the plan needs {steps} steps but the budget allows {allowed}")]
TooManySteps { steps: usize, allowed: usize },
}