use std::collections::{BTreeMap, BTreeSet, VecDeque};
use etdl_parser::ast::{EtlDocument, Node};
use crate::validate::Diagnostic;
const SAFETY_SUPPLEMENT: &str = "etdl.safety";
pub const SAFETY_SCHEMA: &str = "etdl.safety/1.0";
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum Severity {
Catastrophic,
Critical,
Marginal,
Negligible,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum Likelihood {
Frequent,
Probable,
Occasional,
Remote,
Improbable,
}
fn parse_severity(s: &str) -> Option<Severity> {
match s {
"catastrophic" => Some(Severity::Catastrophic),
"critical" => Some(Severity::Critical),
"marginal" => Some(Severity::Marginal),
"negligible" => Some(Severity::Negligible),
_ => None,
}
}
fn parse_likelihood(s: &str) -> Option<Likelihood> {
match s {
"frequent" => Some(Likelihood::Frequent),
"probable" => Some(Likelihood::Probable),
"occasional" => Some(Likelihood::Occasional),
"remote" => Some(Likelihood::Remote),
"improbable" => Some(Likelihood::Improbable),
_ => None,
}
}
fn risk_matrix_value(severity: Severity, likelihood: Likelihood) -> i64 {
use Likelihood::*;
use Severity::*;
match (severity, likelihood) {
(Catastrophic, Frequent) => 1,
(Catastrophic, Probable) => 1,
(Catastrophic, Occasional) => 1,
(Catastrophic, Remote) => 2,
(Catastrophic, Improbable) => 2,
(Critical, Frequent) => 1,
(Critical, Probable) => 1,
(Critical, Occasional) => 2,
(Critical, Remote) => 2,
(Critical, Improbable) => 3,
(Marginal, Frequent) => 1,
(Marginal, Probable) => 2,
(Marginal, Occasional) => 3,
(Marginal, Remote) => 3,
(Marginal, Improbable) => 4,
(Negligible, Frequent) => 2,
(Negligible, Probable) => 3,
(Negligible, Occasional) => 4,
(Negligible, Remote) => 4,
(Negligible, Improbable) => 4,
}
}
#[derive(Debug, Clone, serde::Deserialize)]
pub struct Hazard {
pub id: String,
pub description: String,
pub severity: String,
pub likelihood: String,
#[serde(rename = "riskIndex")]
pub risk_index: i64,
#[serde(rename = "consequenceRef")]
pub consequence_ref: String,
}
#[derive(Debug, Clone, serde::Deserialize)]
pub struct SafetyBarrier {
pub id: String,
#[serde(rename = "nodeRef")]
pub node_ref: String,
pub sil: i64,
#[serde(default, rename = "independentOf")]
pub independent_of: Vec<String>,
#[serde(default, rename = "commonCauseGroup")]
pub common_cause_group: Option<String>,
}
#[derive(Debug, Clone, Default)]
pub struct SafetyData {
pub hazards: Vec<Hazard>,
pub barriers: Vec<SafetyBarrier>,
}
pub fn parse_and_validate_safety(doc: &EtlDocument) -> (SafetyData, Vec<Diagnostic>) {
let mut diagnostics = Vec::new();
let mut data = SafetyData::default();
if !crate::validate::declares_supplement(doc, SAFETY_SUPPLEMENT) {
return (data, diagnostics);
}
let Some(ext) = doc.extensions.get("x-safety") else {
return (data, diagnostics);
};
if let Some(raw_hazards) = ext.get("hazards") {
match serde_yaml::from_value::<Vec<Hazard>>(raw_hazards.clone()) {
Ok(candidates) => {
let mut seen_ids = BTreeSet::new();
for hazard in candidates {
let mut has_error = false;
if !seen_ids.insert(hazard.id.clone()) {
diagnostics.push(Diagnostic::error(
"E-130",
format!("x-safety: duplicate hazard id '{}'", hazard.id),
));
has_error = true;
}
let severity = parse_severity(&hazard.severity);
if severity.is_none() {
diagnostics.push(Diagnostic::error(
"E-130",
format!(
"x-safety: hazard '{}': severity '{}' is not one of catastrophic, critical, marginal, negligible",
hazard.id, hazard.severity
),
));
has_error = true;
}
let likelihood = parse_likelihood(&hazard.likelihood);
if likelihood.is_none() {
diagnostics.push(Diagnostic::error(
"E-130",
format!(
"x-safety: hazard '{}': likelihood '{}' is not one of frequent, probable, occasional, remote, improbable",
hazard.id, hazard.likelihood
),
));
has_error = true;
}
let risk_index_valid = (1..=4).contains(&hazard.risk_index);
if !risk_index_valid {
diagnostics.push(Diagnostic::error(
"E-130",
format!(
"x-safety: hazard '{}': riskIndex must be an integer in [1,4] (got {})",
hazard.id, hazard.risk_index
),
));
has_error = true;
}
if !resolve_node_of_kind(doc, &hazard.consequence_ref, |n| {
matches!(n, Node::Consequence(_))
}) {
diagnostics.push(Diagnostic::error(
"E-131",
format!(
"x-safety: hazard '{}': consequenceRef '{}' does not resolve to a Consequence node",
hazard.id, hazard.consequence_ref
),
));
has_error = true;
}
if let (Some(sev), Some(like), true) = (severity, likelihood, risk_index_valid) {
let expected = risk_matrix_value(sev, like);
if hazard.risk_index != expected {
diagnostics.push(Diagnostic::warning(
"W-410",
format!(
"x-safety: hazard '{}': declared riskIndex {} does not match the risk matrix value {} for severity='{}'/likelihood='{}'",
hazard.id, hazard.risk_index, expected, hazard.severity, hazard.likelihood
),
));
}
}
if !has_error {
data.hazards.push(hazard);
}
}
}
Err(e) => {
diagnostics.push(Diagnostic::error(
"E-130",
format!("x-safety: invalid hazard manifest: {e}"),
));
}
}
}
if let Some(raw_barriers) = ext.get("barriers") {
match serde_yaml::from_value::<Vec<SafetyBarrier>>(raw_barriers.clone()) {
Ok(candidates) => {
let mut seen_ids = BTreeSet::new();
let mut valid_barriers = Vec::new();
for barrier in candidates {
let mut has_error = false;
if !seen_ids.insert(barrier.id.clone()) {
diagnostics.push(Diagnostic::error(
"E-130",
format!("x-safety: duplicate safety barrier id '{}'", barrier.id),
));
has_error = true;
}
if !(1..=4).contains(&barrier.sil) {
diagnostics.push(Diagnostic::error(
"E-130",
format!(
"x-safety: safety barrier '{}': sil must be an integer in [1,4] (got {})",
barrier.id, barrier.sil
),
));
has_error = true;
}
if !resolve_node_of_kind(doc, &barrier.node_ref, |n| matches!(n, Node::Barrier(_))) {
diagnostics.push(Diagnostic::error(
"E-131",
format!(
"x-safety: safety barrier '{}': nodeRef '{}' does not resolve to a Barrier node",
barrier.id, barrier.node_ref
),
));
has_error = true;
}
if !has_error {
valid_barriers.push(barrier);
}
}
check_common_cause_contradictions(&valid_barriers, &mut diagnostics);
data.barriers = valid_barriers;
}
Err(e) => {
diagnostics.push(Diagnostic::error(
"E-130",
format!("x-safety: invalid safety barrier manifest: {e}"),
));
}
}
}
(data, diagnostics)
}
fn resolve_node_of_kind(doc: &EtlDocument, node_ref: &str, is_kind: impl Fn(&Node) -> bool) -> bool {
let rest = node_ref.trim_start_matches('#');
let Some(after) = rest.strip_prefix("/eventTrees/") else {
return false;
};
match after.split('/').collect::<Vec<_>>().as_slice() {
[tree_id, "nodes", node_id] if !tree_id.is_empty() && !node_id.is_empty() => doc
.event_trees
.get(*tree_id)
.and_then(|t| t.nodes.get(*node_id))
.is_some_and(is_kind),
_ => false,
}
}
fn check_common_cause_contradictions(barriers: &[SafetyBarrier], diagnostics: &mut Vec<Diagnostic>) {
let by_id: BTreeMap<&str, &SafetyBarrier> = barriers.iter().map(|b| (b.id.as_str(), b)).collect();
let mut adjacency: BTreeMap<&str, BTreeSet<&str>> = BTreeMap::new();
for barrier in barriers {
for other_id in &barrier.independent_of {
if let Some(other) = by_id.get(other_id.as_str()) {
let mutual = other.independent_of.iter().any(|id| id == &barrier.id);
if mutual {
adjacency.entry(barrier.id.as_str()).or_default().insert(other.id.as_str());
adjacency.entry(other.id.as_str()).or_default().insert(barrier.id.as_str());
}
}
}
}
let mut visited: BTreeSet<&str> = BTreeSet::new();
for barrier in barriers {
if visited.contains(barrier.id.as_str()) {
continue;
}
let mut component = Vec::new();
let mut queue = VecDeque::new();
queue.push_back(barrier.id.as_str());
visited.insert(barrier.id.as_str());
while let Some(id) = queue.pop_front() {
component.push(id);
if let Some(neighbors) = adjacency.get(id) {
for &n in neighbors {
if visited.insert(n) {
queue.push_back(n);
}
}
}
}
let mut groups: BTreeMap<&str, Vec<&str>> = BTreeMap::new();
for id in &component {
if let Some(group) = by_id[id].common_cause_group.as_deref() {
if !group.is_empty() {
groups.entry(group).or_default().push(id);
}
}
}
for (group, mut members) in groups {
if members.len() >= 2 {
members.sort_unstable();
diagnostics.push(Diagnostic::error(
"E-132",
format!(
"x-safety: barriers {} mutually claim independentOf each other but share commonCauseGroup '{}' — self-contradictory",
members.join(", "),
group
),
));
}
}
}
}
#[derive(Debug, Default)]
pub struct SafetyExtension;
impl SafetyExtension {
pub fn new() -> Self {
SafetyExtension
}
}
pub struct SafetyResult {
pub hazards: Vec<Hazard>,
pub barriers: Vec<SafetyBarrier>,
}
impl crate::extension::ExtensionResult for SafetyResult {
fn extension_id(&self) -> &str {
SAFETY_SUPPLEMENT
}
}
impl crate::extension::EtdlExtension for SafetyExtension {
fn id(&self) -> &str {
SAFETY_SUPPLEMENT
}
fn version(&self) -> &str {
"1.0"
}
fn descriptor(&self) -> crate::extension::SupplementDescriptor {
crate::extension::SupplementDescriptor {
summary: "Hazard classification against a severity/likelihood risk matrix, and \
Safety Integrity Level/independence declarations on core Barrier nodes; \
no new probability mathematics.",
schema: Some(SAFETY_SCHEMA),
diagnostic_codes: &["E-130", "E-131", "E-132", "W-410"],
requires: &[],
}
}
fn validate(
&self,
doc: &EtlDocument,
_context: &crate::extension::ExtensionContext<'_>,
diagnostics: &mut Vec<Diagnostic>,
) {
let (_data, safety_diagnostics) = parse_and_validate_safety(doc);
diagnostics.extend(safety_diagnostics);
}
fn process(
&self,
doc: &EtlDocument,
_context: &crate::extension::ExtensionContext<'_>,
_diagnostics: &mut Vec<Diagnostic>,
) -> Box<dyn crate::extension::ExtensionResult + '_> {
let (data, _safety_diagnostics) = parse_and_validate_safety(doc);
Box::new(SafetyResult {
hazards: data.hazards,
barriers: data.barriers,
})
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::extension::{builtin_registry, EtdlExtension, ExtensionContext};
fn doc_with_safety(x_safety_yaml: &str) -> EtlDocument {
let yaml = format!(
r##"
etdl: "1.0.0"
info: {{ title: "T", version: "1.0.0", domain: "D" }}
supplements:
- id: etdl.safety
version: "1.0"
eventTrees:
OrderFulfillment:
initiatingEvent: {{ id: I, message: "a#/m", next: RetryBarrier }}
nodes:
RetryBarrier:
type: barrier
branches:
- outcome: SUCCESS
condition: default
probability: 1.0
next: ProcessPaymentOperation
ProcessPaymentOperation:
type: operation
action: execute
handler: "h"
next: FulfillmentConsequence
FulfillmentConsequence:
type: consequence
operation: terminate
x-safety:
{x_safety_yaml}
"##
);
serde_yaml::from_str(&yaml).unwrap()
}
#[test]
fn safety_extension_is_registered_and_built_in() {
let registry = builtin_registry();
assert!(registry.contains(SAFETY_SUPPLEMENT));
assert!(registry.list().contains(&SAFETY_SUPPLEMENT));
}
#[test]
fn document_without_x_safety_has_no_diagnostics() {
let yaml = r#"
etdl: "1.0.0"
info: { title: "T", version: "1.0.0", domain: "D" }
eventTrees:
T:
initiatingEvent: { id: I, message: "a#/m", next: C }
nodes:
C: { type: consequence, operation: terminate }
"#;
let doc: EtlDocument = serde_yaml::from_str(yaml).unwrap();
let (data, diagnostics) = parse_and_validate_safety(&doc);
assert!(data.hazards.is_empty());
assert!(data.barriers.is_empty());
assert!(diagnostics.is_empty());
}
#[test]
fn valid_hazard_and_barrier_have_no_diagnostics() {
let doc = doc_with_safety(
r##" hazards:
- id: gateway-unavailable
description: "payment cannot be captured while the gateway is down"
severity: critical
likelihood: remote
riskIndex: 2
consequenceRef: "#/eventTrees/OrderFulfillment/nodes/FulfillmentConsequence"
barriers:
- id: retry-barrier
nodeRef: "#/eventTrees/OrderFulfillment/nodes/RetryBarrier"
sil: 2
"##,
);
let (data, diagnostics) = parse_and_validate_safety(&doc);
assert!(diagnostics.is_empty(), "unexpected: {diagnostics:?}");
assert_eq!(data.hazards.len(), 1);
assert_eq!(data.barriers.len(), 1);
}
#[test]
fn missing_hazards_and_barriers_keys_are_not_an_error() {
let doc = doc_with_safety(" {}");
let (data, diagnostics) = parse_and_validate_safety(&doc);
assert!(data.hazards.is_empty());
assert!(data.barriers.is_empty());
assert!(diagnostics.is_empty());
}
#[test]
fn malformed_hazards_produces_e130() {
let doc = doc_with_safety(" hazards: \"oops\"");
let (data, diagnostics) = parse_and_validate_safety(&doc);
assert!(data.hazards.is_empty());
assert!(diagnostics.iter().any(|d| d.code == "E-130"));
}
#[test]
fn invalid_severity_produces_e130() {
let doc = doc_with_safety(
r##" hazards:
- id: h1
description: "d"
severity: extremely-bad
likelihood: remote
riskIndex: 2
consequenceRef: "#/eventTrees/OrderFulfillment/nodes/FulfillmentConsequence"
"##,
);
let (data, diagnostics) = parse_and_validate_safety(&doc);
assert!(data.hazards.is_empty());
assert!(diagnostics.iter().any(|d| d.code == "E-130"));
}
#[test]
fn invalid_likelihood_produces_e130() {
let doc = doc_with_safety(
r##" hazards:
- id: h1
description: "d"
severity: critical
likelihood: all-the-time
riskIndex: 2
consequenceRef: "#/eventTrees/OrderFulfillment/nodes/FulfillmentConsequence"
"##,
);
let (_data, diagnostics) = parse_and_validate_safety(&doc);
assert!(diagnostics.iter().any(|d| d.code == "E-130"));
}
#[test]
fn out_of_range_risk_index_produces_e130() {
let doc = doc_with_safety(
r##" hazards:
- id: h1
description: "d"
severity: critical
likelihood: remote
riskIndex: 9
consequenceRef: "#/eventTrees/OrderFulfillment/nodes/FulfillmentConsequence"
"##,
);
let (_data, diagnostics) = parse_and_validate_safety(&doc);
assert!(diagnostics.iter().any(|d| d.code == "E-130"));
assert!(!diagnostics.iter().any(|d| d.code == "W-410"));
}
#[test]
fn consequence_ref_at_wrong_node_kind_produces_e131() {
let doc = doc_with_safety(
r##" hazards:
- id: h1
description: "d"
severity: critical
likelihood: remote
riskIndex: 2
consequenceRef: "#/eventTrees/OrderFulfillment/nodes/RetryBarrier"
"##,
);
let (data, diagnostics) = parse_and_validate_safety(&doc);
assert!(data.hazards.is_empty());
assert!(diagnostics.iter().any(|d| d.code == "E-131"));
}
#[test]
fn barrier_node_ref_at_wrong_node_kind_produces_e131() {
let doc = doc_with_safety(
r##" barriers:
- id: b1
nodeRef: "#/eventTrees/OrderFulfillment/nodes/ProcessPaymentOperation"
sil: 2
"##,
);
let (data, diagnostics) = parse_and_validate_safety(&doc);
assert!(data.barriers.is_empty());
assert!(diagnostics.iter().any(|d| d.code == "E-131"));
}
#[test]
fn unresolvable_node_ref_produces_e131() {
let doc = doc_with_safety(
r##" barriers:
- id: b1
nodeRef: "#/eventTrees/OrderFulfillment/nodes/DoesNotExist"
sil: 2
"##,
);
let (_data, diagnostics) = parse_and_validate_safety(&doc);
assert!(diagnostics.iter().any(|d| d.code == "E-131"));
}
#[test]
fn out_of_range_sil_produces_e130() {
let doc = doc_with_safety(
r##" barriers:
- id: b1
nodeRef: "#/eventTrees/OrderFulfillment/nodes/RetryBarrier"
sil: 7
"##,
);
let (data, diagnostics) = parse_and_validate_safety(&doc);
assert!(data.barriers.is_empty());
assert!(diagnostics.iter().any(|d| d.code == "E-130"));
}
#[test]
fn duplicate_hazard_id_produces_e130() {
let doc = doc_with_safety(
r##" hazards:
- id: dup
description: "d1"
severity: critical
likelihood: remote
riskIndex: 2
consequenceRef: "#/eventTrees/OrderFulfillment/nodes/FulfillmentConsequence"
- id: dup
description: "d2"
severity: marginal
likelihood: probable
riskIndex: 2
consequenceRef: "#/eventTrees/OrderFulfillment/nodes/FulfillmentConsequence"
"##,
);
let (_data, diagnostics) = parse_and_validate_safety(&doc);
assert!(diagnostics.iter().any(|d| d.code == "E-130" && d.message.contains("duplicate hazard id")));
}
#[test]
fn duplicate_barrier_id_produces_e130() {
let doc = doc_with_safety(
r##" barriers:
- id: dup
nodeRef: "#/eventTrees/OrderFulfillment/nodes/RetryBarrier"
sil: 1
- id: dup
nodeRef: "#/eventTrees/OrderFulfillment/nodes/RetryBarrier"
sil: 2
"##,
);
let (_data, diagnostics) = parse_and_validate_safety(&doc);
assert!(diagnostics.iter().any(|d| d.code == "E-130" && d.message.contains("duplicate safety barrier id")));
}
#[test]
fn mismatched_risk_index_produces_w410() {
let doc = doc_with_safety(
r##" hazards:
- id: h1
description: "d"
severity: catastrophic
likelihood: remote
riskIndex: 4
consequenceRef: "#/eventTrees/OrderFulfillment/nodes/FulfillmentConsequence"
"##,
);
let (data, diagnostics) = parse_and_validate_safety(&doc);
assert_eq!(data.hazards.len(), 1);
assert!(diagnostics.iter().any(|d| d.code == "W-410"));
assert!(!diagnostics.iter().any(|d| d.is_error()));
}
#[test]
fn matching_risk_index_has_no_w410() {
let doc = doc_with_safety(
r##" hazards:
- id: h1
description: "d"
severity: catastrophic
likelihood: remote
riskIndex: 2
consequenceRef: "#/eventTrees/OrderFulfillment/nodes/FulfillmentConsequence"
"##,
);
let (_data, diagnostics) = parse_and_validate_safety(&doc);
assert!(diagnostics.is_empty(), "unexpected: {diagnostics:?}");
}
#[test]
fn mutual_independent_of_with_shared_common_cause_group_produces_e132() {
let doc = doc_with_safety(
r##" barriers:
- id: retry-barrier
nodeRef: "#/eventTrees/OrderFulfillment/nodes/RetryBarrier"
sil: 2
independentOf: ["fallback-barrier"]
commonCauseGroup: "primary-network-path"
- id: fallback-barrier
nodeRef: "#/eventTrees/OrderFulfillment/nodes/RetryBarrier"
sil: 1
independentOf: ["retry-barrier"]
commonCauseGroup: "primary-network-path"
"##,
);
let (data, diagnostics) = parse_and_validate_safety(&doc);
assert_eq!(data.barriers.len(), 2);
assert!(diagnostics.iter().any(|d| d.code == "E-132"));
}
#[test]
fn one_sided_independent_of_does_not_produce_e132() {
let doc = doc_with_safety(
r##" barriers:
- id: retry-barrier
nodeRef: "#/eventTrees/OrderFulfillment/nodes/RetryBarrier"
sil: 2
independentOf: ["fallback-barrier"]
commonCauseGroup: "primary-network-path"
- id: fallback-barrier
nodeRef: "#/eventTrees/OrderFulfillment/nodes/RetryBarrier"
sil: 1
commonCauseGroup: "primary-network-path"
"##,
);
let (_data, diagnostics) = parse_and_validate_safety(&doc);
assert!(!diagnostics.iter().any(|d| d.code == "E-132"));
}
#[test]
fn mutual_independent_of_with_different_common_cause_groups_does_not_produce_e132() {
let doc = doc_with_safety(
r##" barriers:
- id: retry-barrier
nodeRef: "#/eventTrees/OrderFulfillment/nodes/RetryBarrier"
sil: 2
independentOf: ["fallback-barrier"]
commonCauseGroup: "primary-network-path"
- id: fallback-barrier
nodeRef: "#/eventTrees/OrderFulfillment/nodes/RetryBarrier"
sil: 1
independentOf: ["retry-barrier"]
commonCauseGroup: "secondary-network-path"
"##,
);
let (_data, diagnostics) = parse_and_validate_safety(&doc);
assert!(!diagnostics.iter().any(|d| d.code == "E-132"));
}
#[test]
fn process_returns_typed_result_with_correct_extension_id() {
let doc = doc_with_safety(
r##" barriers:
- id: retry-barrier
nodeRef: "#/eventTrees/OrderFulfillment/nodes/RetryBarrier"
sil: 2
"##,
);
let ext = SafetyExtension::new();
let base = std::path::Path::new(".");
let ctx = ExtensionContext::new(&doc, base);
let mut diagnostics = Vec::new();
let result = ext.process(&doc, &ctx, &mut diagnostics);
assert!(diagnostics.is_empty(), "unexpected: {diagnostics:?}");
assert_eq!(result.extension_id(), SAFETY_SUPPLEMENT);
assert!(result.basic_event_overrides().is_empty());
}
}