use indexmap::IndexMap;
use serde::{Deserialize, Serialize};
use std::collections::BTreeMap;
use thiserror::Error;
pub const DEFAULT_PORT_BASE: u16 = 34500;
pub const LOCAL_ADDRESS: &str = "127.0.0.1";
pub const ENV_PARTICIPANTS: &str = "QED_PARTICIPANTS";
pub const ENV_PARTICIPANT_SELF: &str = "QED_PARTICIPANT_SELF";
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
#[cfg_attr(feature = "json-schema", derive(schemars::JsonSchema))]
#[serde(deny_unknown_fields)]
pub struct ParticipantSpec {
#[serde(default, skip_serializing_if = "Option::is_none")]
pub node: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub address: Option<String>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub ports: Vec<String>,
#[serde(default, skip_serializing_if = "std::ops::Not::not")]
pub coordinator: bool,
}
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
#[cfg_attr(feature = "json-schema", derive(schemars::JsonSchema))]
#[serde(deny_unknown_fields)]
pub struct ParticipantSet {
#[serde(default, skip_serializing_if = "Option::is_none")]
pub port_base: Option<u16>,
#[serde(default, rename = "role")]
#[cfg_attr(
feature = "json-schema",
schemars(schema_with = "crate::types::permissive_schema")
)]
pub roles: IndexMap<String, ParticipantSpec>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct Participant {
pub name: String,
#[serde(skip_serializing_if = "Option::is_none")]
pub node: Option<String>,
pub address: String,
pub ports: BTreeMap<String, u16>,
pub coordinator: bool,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ParticipantPlan {
participants: Vec<Participant>,
}
#[derive(Debug, Clone, PartialEq, Eq, Error)]
pub enum ParticipantError {
#[error("[pipeline.participants] declares no roles — drop the block or declare one")]
Empty,
#[error(
"[pipeline.participants] declares no coordinator — exactly one role needs \
`coordinator = true`, and its exit status is the run's verdict"
)]
NoCoordinator,
#[error(
"[pipeline.participants] declares {count} coordinators ({names}) — exactly \
one role may be the coordinator, because a run has one verdict"
)]
ManyCoordinators { count: usize, names: String },
#[error(
"participant `{role}`: role names must be non-empty and use only \
[A-Za-z0-9_-] (they travel as ids in {env})",
env = ENV_PARTICIPANT_SELF
)]
BadRoleName { role: String },
#[error(
"participant `{role}`: declares `node = \"{node}\"` but no `address` — QED \
has no machine inventory, and a rendezvous cannot be built from a name \
its peers cannot dial. Add `address = \"<ip-or-host>\"`."
)]
RemoteWithoutAddress { role: String, node: String },
#[error("participant `{role}`: port name `{port}` is declared twice")]
DuplicatePort { role: String, port: String },
#[error(
"participant `{role}`: port names must be non-empty and use only \
[A-Za-z0-9_-]"
)]
BadPortName { role: String, port: String },
#[error(
"[pipeline.participants] needs {needed} ports from port_base {base}, which \
runs past 65535 — lower `port_base` or declare fewer ports"
)]
PortRangeExhausted { needed: usize, base: u16 },
#[error(
"step `{step}`: `participant = \"{participant}\"` names no declared role — \
[pipeline.participants] declares {declared}"
)]
UnknownRole {
step: String,
participant: String,
declared: String,
},
#[error(
"step `{step}`: `participant = \"{participant}\"` but the pipeline declares no \
[pipeline.participants] block — the role has no host, no address and no \
place in the verdict"
)]
StepWithoutSet { step: String, participant: String },
#[error(
"participant `{role}`: no step declares `participant = \"{role}\"` — a role \
with no steps can only ever report as never-started, so this is a typo, \
not a configuration"
)]
RoleWithoutSteps { role: String },
}
fn is_token(s: &str) -> bool {
!s.is_empty()
&& s.chars()
.all(|c| c.is_ascii_alphanumeric() || c == '_' || c == '-')
}
impl ParticipantSet {
pub fn is_empty(&self) -> bool {
self.roles.is_empty()
}
pub fn plan(&self) -> Result<ParticipantPlan, ParticipantError> {
if self.roles.is_empty() {
return Err(ParticipantError::Empty);
}
let coordinators: Vec<&String> = self
.roles
.iter()
.filter(|(_, spec)| spec.coordinator)
.map(|(name, _)| name)
.collect();
match coordinators.len() {
0 => return Err(ParticipantError::NoCoordinator),
1 => {}
n => {
return Err(ParticipantError::ManyCoordinators {
count: n,
names: coordinators
.iter()
.map(|s| s.as_str())
.collect::<Vec<_>>()
.join(", "),
})
}
}
let base = self.port_base.unwrap_or(DEFAULT_PORT_BASE);
let needed: usize = self.roles.values().map(|s| s.ports.len()).sum();
if needed > 0 && (base as usize) + needed - 1 > u16::MAX as usize {
return Err(ParticipantError::PortRangeExhausted { needed, base });
}
let mut assigned: usize = 0;
let mut participants = Vec::with_capacity(self.roles.len());
for (name, spec) in &self.roles {
if !is_token(name) {
return Err(ParticipantError::BadRoleName { role: name.clone() });
}
let address = match (&spec.node, &spec.address) {
(_, Some(addr)) => addr.clone(),
(None, None) => LOCAL_ADDRESS.to_string(),
(Some(node), None) => {
return Err(ParticipantError::RemoteWithoutAddress {
role: name.clone(),
node: node.clone(),
})
}
};
let mut ports = BTreeMap::new();
for port in &spec.ports {
if !is_token(port) {
return Err(ParticipantError::BadPortName {
role: name.clone(),
port: port.clone(),
});
}
if ports.contains_key(port) {
return Err(ParticipantError::DuplicatePort {
role: name.clone(),
port: port.clone(),
});
}
ports.insert(port.clone(), (base as usize + assigned) as u16);
assigned += 1;
}
participants.push(Participant {
name: name.clone(),
node: spec.node.clone(),
address,
ports,
coordinator: spec.coordinator,
});
}
Ok(ParticipantPlan { participants })
}
}
impl ParticipantPlan {
pub fn participants(&self) -> &[Participant] {
&self.participants
}
pub fn get(&self, name: &str) -> Option<&Participant> {
self.participants.iter().find(|p| p.name == name)
}
pub fn coordinator(&self) -> &Participant {
self.participants
.iter()
.find(|p| p.coordinator)
.expect("plan() rejects a set without exactly one coordinator")
}
pub fn rendezvous_env(&self, self_name: Option<&str>) -> Vec<(String, String)> {
let json = serde_json::to_string(&self.participants)
.expect("Participant is a plain serde struct with no non-string map keys");
let mut env = vec![(ENV_PARTICIPANTS.to_string(), json)];
if let Some(name) = self_name {
env.push((ENV_PARTICIPANT_SELF.to_string(), name.to_string()));
}
env
}
pub fn teardown_order(&self) -> impl Iterator<Item = &Participant> {
self.participants
.iter()
.filter(|p| !p.coordinator)
.chain(self.participants.iter().filter(|p| p.coordinator))
}
}
pub fn plan_for(pipeline: &crate::types::Pipeline) -> Result<Option<ParticipantPlan>, ParticipantError> {
let all_steps = || pipeline.steps.iter().chain(pipeline.finally.iter());
let Some(set) = pipeline.participants.as_ref().filter(|s| !s.is_empty()) else {
if let Some(step) = all_steps().find(|s| s.participant.is_some()) {
return Err(ParticipantError::StepWithoutSet {
step: step.name.clone(),
participant: step.participant.clone().unwrap_or_default(),
});
}
return Ok(None);
};
let plan = set.plan()?;
for step in all_steps() {
let Some(role) = step.participant.as_deref() else {
continue;
};
if plan.get(role).is_none() {
return Err(ParticipantError::UnknownRole {
step: step.name.clone(),
participant: role.to_string(),
declared: plan
.participants()
.iter()
.map(|p| p.name.as_str())
.collect::<Vec<_>>()
.join(", "),
});
}
}
for participant in plan.participants() {
if !all_steps().any(|s| s.participant.as_deref() == Some(participant.name.as_str())) {
return Err(ParticipantError::RoleWithoutSteps {
role: participant.name.clone(),
});
}
}
Ok(Some(plan))
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum ParticipantOutcome {
NeverStarted { reason: String },
Failed { detail: String },
Completed,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ParticipantReport {
pub name: String,
pub coordinator: bool,
pub outcome: ParticipantOutcome,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum Verdict {
Pass,
Fail { summary: String },
}
impl Verdict {
pub fn passed(&self) -> bool {
matches!(self, Verdict::Pass)
}
}
pub fn verdict(reports: &[ParticipantReport]) -> Verdict {
let never: Vec<String> = reports
.iter()
.filter_map(|r| match &r.outcome {
ParticipantOutcome::NeverStarted { reason } => {
Some(format!("`{}` never started: {reason}", r.name))
}
_ => None,
})
.collect();
if !never.is_empty() {
return Verdict::Fail {
summary: format!(
"{} of {} participants never started — this is a fleet fault, not a \
test failure: {}",
never.len(),
reports.len(),
never.join("; "),
),
};
}
let mut failed: Vec<&ParticipantReport> = reports
.iter()
.filter(|r| matches!(r.outcome, ParticipantOutcome::Failed { .. }))
.collect();
failed.sort_by_key(|r| !r.coordinator);
if !failed.is_empty() {
let detail = failed
.iter()
.map(|r| match &r.outcome {
ParticipantOutcome::Failed { detail } => format!("`{}`: {detail}", r.name),
_ => unreachable!("filtered to Failed"),
})
.collect::<Vec<_>>()
.join("; ");
return Verdict::Fail {
summary: format!("{} participant(s) failed — {detail}", failed.len()),
};
}
if !reports.iter().any(|r| r.coordinator) {
return Verdict::Fail {
summary: "no coordinator report — the participant holding the run's \
verdict left no account of itself"
.to_string(),
};
}
Verdict::Pass
}
#[cfg(test)]
mod tests {
use super::*;
fn set(toml_src: &str) -> ParticipantSet {
toml::from_str(toml_src).expect("participant set parses")
}
#[test]
fn parses_the_documented_shape() {
let s = set(r#"
port_base = 34500
[role.runner]
coordinator = true
[role.responder]
node = "us-west-011"
address = "100.64.0.11"
ports = ["clock"]
"#);
assert_eq!(s.port_base, Some(34500));
assert_eq!(s.roles.len(), 2);
assert_eq!(
s.roles.keys().collect::<Vec<_>>(),
vec!["runner", "responder"]
);
assert!(s.roles["runner"].coordinator);
assert_eq!(s.roles["responder"].node.as_deref(), Some("us-west-011"));
}
#[test]
fn plan_binds_hosts_addresses_and_ports() {
let plan = set(r#"
[role.runner]
coordinator = true
ports = ["control"]
[role.responder]
node = "us-west-011"
address = "100.64.0.11"
ports = ["clock", "aux"]
"#)
.plan()
.unwrap();
let runner = plan.get("runner").unwrap();
assert_eq!(runner.node, None);
assert_eq!(runner.address, LOCAL_ADDRESS);
assert_eq!(runner.ports["control"], DEFAULT_PORT_BASE);
let responder = plan.get("responder").unwrap();
assert_eq!(responder.node.as_deref(), Some("us-west-011"));
assert_eq!(responder.address, "100.64.0.11");
assert_eq!(responder.ports["clock"], DEFAULT_PORT_BASE + 1);
assert_eq!(responder.ports["aux"], DEFAULT_PORT_BASE + 2);
}
#[test]
fn every_assigned_port_is_distinct_across_the_whole_set() {
let plan = set(r#"
[role.a]
coordinator = true
ports = ["x", "y"]
[role.b]
node = "n1"
address = "10.0.0.1"
ports = ["x", "y"]
[role.c]
node = "n2"
address = "10.0.0.2"
ports = ["x"]
"#)
.plan()
.unwrap();
let mut all: Vec<u16> = plan
.participants()
.iter()
.flat_map(|p| p.ports.values().copied())
.collect();
assert_eq!(all.len(), 5);
all.sort_unstable();
all.dedup();
assert_eq!(all.len(), 5, "same port handed to two participants");
}
#[test]
fn plan_is_deterministic() {
let src = r#"
[role.a]
coordinator = true
ports = ["x"]
[role.b]
node = "n1"
address = "10.0.0.1"
ports = ["y"]
"#;
assert_eq!(set(src).plan().unwrap(), set(src).plan().unwrap());
}
#[test]
fn a_named_node_without_an_address_is_refused() {
let err = set(r#"
[role.runner]
coordinator = true
[role.responder]
node = "us-west-011"
"#)
.plan()
.unwrap_err();
assert_eq!(
err,
ParticipantError::RemoteWithoutAddress {
role: "responder".into(),
node: "us-west-011".into(),
}
);
}
#[test]
fn exactly_one_coordinator() {
let none = set(r#"
[role.a]
[role.b]
"#)
.plan()
.unwrap_err();
assert_eq!(none, ParticipantError::NoCoordinator);
let two = set(r#"
[role.a]
coordinator = true
[role.b]
coordinator = true
"#)
.plan()
.unwrap_err();
assert!(matches!(
two,
ParticipantError::ManyCoordinators { count: 2, .. }
));
}
#[test]
fn empty_set_is_an_error_not_a_silent_no_op() {
assert_eq!(
ParticipantSet::default().plan().unwrap_err(),
ParticipantError::Empty
);
assert!(ParticipantSet::default().is_empty());
}
#[test]
fn bad_names_are_refused() {
let mut s = ParticipantSet::default();
s.roles.insert(
"has space".into(),
ParticipantSpec {
coordinator: true,
..Default::default()
},
);
assert!(matches!(
s.plan().unwrap_err(),
ParticipantError::BadRoleName { .. }
));
let bad_port = set(r#"
[role.a]
coordinator = true
ports = ["has space"]
"#)
.plan()
.unwrap_err();
assert!(matches!(bad_port, ParticipantError::BadPortName { .. }));
}
#[test]
fn duplicate_port_name_within_a_role_is_refused() {
let err = set(r#"
[role.a]
coordinator = true
ports = ["clock", "clock"]
"#)
.plan()
.unwrap_err();
assert_eq!(
err,
ParticipantError::DuplicatePort {
role: "a".into(),
port: "clock".into(),
}
);
}
#[test]
fn port_range_exhaustion_is_caught_at_plan_time() {
let err = set(r#"
port_base = 65535
[role.a]
coordinator = true
ports = ["x", "y"]
"#)
.plan()
.unwrap_err();
assert_eq!(
err,
ParticipantError::PortRangeExhausted {
needed: 2,
base: 65535
}
);
assert!(set(r#"
port_base = 65535
[role.a]
coordinator = true
ports = ["x"]
"#)
.plan()
.is_ok());
}
#[test]
fn unknown_keys_are_refused_rather_than_silently_ignored() {
assert!(toml::from_str::<ParticipantSet>(
r#"
[role.a]
coordinator = true
adress = "10.0.0.1"
"#
)
.is_err());
}
#[test]
fn rendezvous_env_carries_the_whole_set_and_who_you_are() {
let plan = set(r#"
[role.runner]
coordinator = true
[role.responder]
node = "us-west-011"
address = "100.64.0.11"
ports = ["clock"]
"#)
.plan()
.unwrap();
let env: BTreeMap<String, String> =
plan.rendezvous_env(Some("responder")).into_iter().collect();
assert_eq!(env[ENV_PARTICIPANT_SELF], "responder");
let decoded: Vec<Participant> = serde_json::from_str(&env[ENV_PARTICIPANTS]).unwrap();
assert_eq!(decoded.len(), 2);
assert_eq!(decoded[0].name, "runner");
assert_eq!(decoded[1].address, "100.64.0.11");
assert_eq!(decoded[1].ports["clock"], DEFAULT_PORT_BASE);
}
#[test]
fn a_step_with_no_participant_still_sees_the_set_but_has_no_self() {
let plan = set(r#"
[role.a]
coordinator = true
"#)
.plan()
.unwrap();
let env: BTreeMap<String, String> = plan.rendezvous_env(None).into_iter().collect();
assert!(env.contains_key(ENV_PARTICIPANTS));
assert!(!env.contains_key(ENV_PARTICIPANT_SELF));
}
#[test]
fn teardown_reaps_peers_first_and_the_coordinator_last() {
let plan = set(r#"
[role.peer_a]
node = "n1"
address = "10.0.0.1"
[role.runner]
coordinator = true
[role.peer_b]
node = "n2"
address = "10.0.0.2"
"#)
.plan()
.unwrap();
let order: Vec<&str> = plan
.teardown_order()
.map(|p| p.name.as_str())
.collect();
assert_eq!(order, vec!["peer_a", "peer_b", "runner"]);
}
fn report(name: &str, coordinator: bool, outcome: ParticipantOutcome) -> ParticipantReport {
ParticipantReport {
name: name.into(),
coordinator,
outcome,
}
}
#[test]
fn all_completed_passes() {
let v = verdict(&[
report("runner", true, ParticipantOutcome::Completed),
report("responder", false, ParticipantOutcome::Completed),
]);
assert_eq!(v, Verdict::Pass);
}
#[test]
fn one_peer_failing_fails_the_run() {
let v = verdict(&[
report("runner", true, ParticipantOutcome::Completed),
report(
"responder",
false,
ParticipantOutcome::Failed {
detail: "bind: address in use".into(),
},
),
]);
let Verdict::Fail { summary } = v else {
panic!("expected a failure")
};
assert!(summary.contains("responder"), "{summary}");
assert!(summary.contains("address in use"), "{summary}");
}
#[test]
fn never_started_outranks_failed_and_says_it_is_a_fleet_fault() {
let v = verdict(&[
report(
"runner",
true,
ParticipantOutcome::Failed {
detail: "connection refused".into(),
},
),
report(
"responder",
false,
ParticipantOutcome::NeverStarted {
reason: "node us-west-011 unreachable".into(),
},
),
]);
let Verdict::Fail { summary } = v else {
panic!("expected a failure")
};
assert!(summary.contains("never started"), "{summary}");
assert!(summary.contains("fleet fault"), "{summary}");
assert!(summary.contains("us-west-011"), "{summary}");
assert!(!summary.contains("connection refused"), "{summary}");
}
#[test]
fn coordinator_failure_is_named_before_a_peers() {
let v = verdict(&[
report(
"responder",
false,
ParticipantOutcome::Failed {
detail: "peer detail".into(),
},
),
report(
"runner",
true,
ParticipantOutcome::Failed {
detail: "assertion detail".into(),
},
),
]);
let Verdict::Fail { summary } = v else {
panic!("expected a failure")
};
let coord_at = summary.find("assertion detail").unwrap();
let peer_at = summary.find("peer detail").unwrap();
assert!(coord_at < peer_at, "{summary}");
}
#[test]
fn a_set_with_no_coordinator_report_does_not_pass() {
let v = verdict(&[report(
"responder",
false,
ParticipantOutcome::Completed,
)]);
assert!(matches!(v, Verdict::Fail { .. }));
}
#[test]
fn an_empty_report_set_does_not_pass() {
assert!(matches!(verdict(&[]), Verdict::Fail { .. }));
}
}