///|
/// Auditable orchestration for turning raw RR data into a usable decision.
/// Every stage has a status, duration supplied by the caller, and explicit
/// reasons so a dashboard can distinguish a failed gate from missing data.
pub(all) enum PipelineStageKind {
PipelineInput
PipelineValidation
PipelineCleaning
PipelineTimeDomain
PipelineFrequency
PipelineNonlinear
PipelineQualityGate
PipelineExport
} derive(FromJson, ToJson, Debug, Eq)
///|
pub(all) enum PipelineStageStatus {
PipelinePending
PipelinePassed
PipelineSkipped
PipelineFailed
} derive(FromJson, ToJson, Debug, Eq)
///|
pub(all) struct PipelineStageTrace {
ordinal : Int
kind : PipelineStageKind
status : PipelineStageStatus
input_count : Int
output_count : Int
quality : Double
duration_ms : Double
message : String
} derive(FromJson, ToJson, Debug, Eq)
///|
pub(all) struct QualityGatePolicy {
minimum_samples : Int
minimum_clean_ratio : Double
minimum_signal_quality : Double
maximum_artifact_ratio : Double
maximum_gap_ratio : Double
require_frequency : Bool
require_nonlinear : Bool
} derive(FromJson, ToJson, Debug, Eq)
///|
pub fn QualityGatePolicy::default() -> QualityGatePolicy {
{
minimum_samples: 32,
minimum_clean_ratio: 0.70,
minimum_signal_quality: 0.70,
maximum_artifact_ratio: 0.30,
maximum_gap_ratio: 0.25,
require_frequency: true,
require_nonlinear: false,
}
}
///|
pub(all) struct QualityGateResult {
passed : Bool
score : Double
clean_ratio : Double
artifact_ratio : Double
gap_ratio : Double
reasons : Array[String]
warnings : Array[String]
} derive(FromJson, ToJson, Debug, Eq)
///|
pub(all) struct HrvPipelineRun {
run_id : String
report : AnalysisReport?
gate : QualityGateResult
stages : Array[PipelineStageTrace]
input_count : Int
accepted : Bool
feature_vector : Array[Double]
export_csv : String
} derive(FromJson, ToJson, Debug, Eq)
///|
pub(all) struct PipelineBatchSummary {
runs : Array[HrvPipelineRun]
total_runs : Int
accepted_runs : Int
rejected_runs : Int
acceptance_ratio : Double
average_quality : Double
average_feature_count : Double
} derive(FromJson, ToJson, Debug, Eq)
///|
pub fn pipeline_stage_name(kind : PipelineStageKind) -> String {
match kind {
PipelineInput => "input"
PipelineValidation => "validation"
PipelineCleaning => "cleaning"
PipelineTimeDomain => "time_domain"
PipelineFrequency => "frequency"
PipelineNonlinear => "nonlinear"
PipelineQualityGate => "quality_gate"
PipelineExport => "export"
}
}
///|
pub fn pipeline_status_name(status : PipelineStageStatus) -> String {
match status {
PipelinePending => "pending"
PipelinePassed => "passed"
PipelineSkipped => "skipped"
PipelineFailed => "failed"
}
}
///|
fn pipeline_bound(value : Double, low : Double, high : Double) -> Double {
if value.is_nan() || value.is_inf() {
low
} else {
value.clamp(min=low, max=high)
}
}
///|
fn pipeline_safe_ratio(numerator : Int, denominator : Int) -> Double {
if denominator <= 0 {
0.0
} else {
numerator.to_double() / denominator.to_double()
}
}
///|
fn pipeline_bool_double(value : Bool) -> Double {
if value {
1.0
} else {
0.0
}
}
///|
pub fn make_quality_gate(
report : AnalysisReport,
policy : QualityGatePolicy,
) -> QualityGateResult {
let total = report.raw_count
let valid = report.quality.valid_beats
let clean_ratio = report.quality.clean_ratio
let artifact_ratio = if total <= 0 {
1.0
} else {
(total - valid).max(0).to_double() / total.to_double()
}
let gaps = report.raw_validation.invalid
let gap_ratio = pipeline_safe_ratio(gaps, total)
let reasons = []
let warnings = []
if total < policy.minimum_samples {
reasons.push("recording has fewer samples than the minimum gate")
}
if clean_ratio < policy.minimum_clean_ratio {
reasons.push("clean interval ratio is below the minimum gate")
}
if report.quality.clean_ratio < policy.minimum_signal_quality {
reasons.push("signal quality is below the preferred floor")
}
if artifact_ratio > policy.maximum_artifact_ratio {
reasons.push("artifact ratio exceeds the configured maximum")
}
if gap_ratio > policy.maximum_gap_ratio {
reasons.push(
"invalid or missing interval ratio exceeds the configured maximum",
)
}
if report.frequency.total_power <= 0.0 && policy.require_frequency {
reasons.push("frequency stage did not produce usable power")
}
if report.nonlinear.complexity_index == 0.0 && policy.require_nonlinear {
reasons.push(
"nonlinear stage was required but returned no complexity index",
)
}
if total > 0 && clean_ratio < 0.85 {
warnings.push("cleaning changed a material part of the input")
}
if report.segments.length() == 0 {
warnings.push("recording is shorter than the configured segment size")
}
let score = (clean_ratio * 0.55 +
(1.0 - artifact_ratio).clamp(min=0.0, max=1.0) * 0.25 +
(1.0 - gap_ratio).clamp(min=0.0, max=1.0) * 0.20).clamp(min=0.0, max=1.0)
{
passed: reasons.length() == 0,
score,
clean_ratio,
artifact_ratio,
gap_ratio,
reasons,
warnings,
}
}
///|
fn pipeline_stage(
ordinal : Int,
kind : PipelineStageKind,
status : PipelineStageStatus,
input_count : Int,
output_count : Int,
quality : Double,
duration_ms : Double,
message : String,
) -> PipelineStageTrace {
{
ordinal,
kind,
status,
input_count,
output_count,
quality: pipeline_bound(quality, 0.0, 1.0),
duration_ms: pipeline_bound(duration_ms, 0.0, 86400000.0),
message,
}
}
///|
pub fn run_hrv_quality_pipeline(
run_id : String,
intervals : Array[Double],
config : HrvConfig,
options : AnalysisOptions,
policy : QualityGatePolicy,
) -> HrvPipelineRun {
let stages = []
stages.push(
pipeline_stage(
0,
PipelineInput,
if intervals.length() == 0 {
PipelineFailed
} else {
PipelinePassed
},
intervals.length(),
intervals.length(),
if intervals.length() == 0 {
0.0
} else {
1.0
},
0.0,
if intervals.length() == 0 {
"empty input"
} else {
"raw RR input accepted"
},
),
)
if intervals.length() == 0 {
let gate = {
passed: false,
score: 0.0,
clean_ratio: 0.0,
artifact_ratio: 1.0,
gap_ratio: 1.0,
reasons: ["input is empty"],
warnings: [],
}
return {
run_id,
report: None,
gate,
stages,
input_count: 0,
accepted: false,
feature_vector: [],
export_csv: "",
}
}
let report = analyze_rr(intervals, config, options)
stages.push(
pipeline_stage(
1,
PipelineValidation,
if report.raw_validation.invalid == 0 {
PipelinePassed
} else {
PipelineFailed
},
intervals.length(),
report.raw_validation.valid,
report.quality.clean_ratio,
0.0,
"interval validation completed",
),
)
stages.push(
pipeline_stage(
2,
PipelineCleaning,
PipelinePassed,
intervals.length(),
report.cleaned_intervals.length(),
report.quality.clean_ratio,
0.0,
"cleaning and artifact replacement completed",
),
)
stages.push(
pipeline_stage(
3,
PipelineTimeDomain,
PipelinePassed,
report.cleaned_intervals.length(),
5,
report.quality.clean_ratio,
0.0,
"time-domain metrics computed",
),
)
stages.push(
pipeline_stage(
4,
PipelineFrequency,
if report.frequency.total_power > 0.0 {
PipelinePassed
} else {
PipelineSkipped
},
report.cleaned_intervals.length(),
7,
report.quality.clean_ratio,
0.0,
"frequency metrics computed",
),
)
stages.push(
pipeline_stage(
5,
PipelineNonlinear,
if options.nonlinear_enabled {
PipelinePassed
} else {
PipelineSkipped
},
report.cleaned_intervals.length(),
8,
report.quality.clean_ratio,
0.0,
if options.nonlinear_enabled {
"nonlinear metrics computed"
} else {
"nonlinear metrics disabled"
},
),
)
let gate = make_quality_gate(report, policy)
stages.push(
pipeline_stage(
6,
PipelineQualityGate,
if gate.passed {
PipelinePassed
} else {
PipelineFailed
},
report.raw_count,
report.cleaned_intervals.length(),
gate.score,
0.0,
if gate.passed {
"quality gate passed"
} else {
"quality gate rejected the run"
},
),
)
let feature_vector = if gate.passed { report.feature_vector } else { [] }
let exported_csv = if gate.passed {
report_table_to_csv(build_report_table(report))
} else {
""
}
stages.push(
pipeline_stage(
7,
PipelineExport,
if gate.passed {
PipelinePassed
} else {
PipelineSkipped
},
report.cleaned_intervals.length(),
feature_vector.length(),
gate.score,
0.0,
if gate.passed {
"accepted feature export generated"
} else {
"export skipped after gate failure"
},
),
)
{
run_id,
report: Some(report),
gate,
stages,
input_count: intervals.length(),
accepted: gate.passed,
feature_vector,
export_csv: exported_csv,
}
}
///|
pub fn pipeline_run_is_usable(run : HrvPipelineRun) -> Bool {
run.accepted &&
run.report is Some(_) &&
run.gate.passed &&
run.feature_vector.length() > 0
}
///|
pub fn pipeline_run_quality(run : HrvPipelineRun) -> Double {
run.gate.score
}
///|
pub fn pipeline_run_stage(
run : HrvPipelineRun,
kind : PipelineStageKind,
) -> PipelineStageTrace? {
for stage in run.stages {
if stage.kind == kind {
return Some(stage)
}
}
None
}
///|
pub fn pipeline_run_failure_reasons(run : HrvPipelineRun) -> Array[String] {
run.gate.reasons
}
///|
pub fn pipeline_run_warning_count(run : HrvPipelineRun) -> Int {
run.gate.warnings.length()
}
///|
pub fn pipeline_run_feature_count(run : HrvPipelineRun) -> Int {
run.feature_vector.length()
}
///|
pub fn pipeline_run_duration(run : HrvPipelineRun) -> Double {
sum_values(run.stages.map(stage => stage.duration_ms))
}
///|
pub fn pipeline_stage_counts(run : HrvPipelineRun) -> Array[Int] {
[
run.stages.filter(stage => stage.status == PipelinePassed).length(),
run.stages.filter(stage => stage.status == PipelineSkipped).length(),
run.stages.filter(stage => stage.status == PipelineFailed).length(),
run.stages.filter(stage => stage.status == PipelinePending).length(),
]
}
///|
pub fn pipeline_batch_summary(
runs : Array[HrvPipelineRun],
) -> PipelineBatchSummary {
let accepted = runs.filter(run => run.accepted).length()
let rejected = runs.length() - accepted
let qualities = runs.map(run => run.gate.score)
let feature_counts = runs.map(run => run.feature_vector.length().to_double())
{
runs,
total_runs: runs.length(),
accepted_runs: accepted,
rejected_runs: rejected,
acceptance_ratio: if runs.length() == 0 {
0.0
} else {
accepted.to_double() / runs.length().to_double()
},
average_quality: mean_value(qualities),
average_feature_count: mean_value(feature_counts),
}
}
///|
pub fn pipeline_batch_is_stable(summary : PipelineBatchSummary) -> Bool {
summary.total_runs > 0 &&
summary.acceptance_ratio >= 0.70 &&
summary.average_quality >= 0.60
}
///|
pub fn pipeline_batch_csv(summary : PipelineBatchSummary) -> String {
let grid = [
[
"run_id", "accepted", "quality_score", "clean_ratio", "artifact_ratio", "gap_ratio",
"feature_count", "failure_count", "warning_count",
],
]
for run in summary.runs {
grid.push([
run.run_id,
run.accepted.to_string(),
run.gate.score.to_string(),
run.gate.clean_ratio.to_string(),
run.gate.artifact_ratio.to_string(),
run.gate.gap_ratio.to_string(),
run.feature_vector.length().to_string(),
run.gate.reasons.length().to_string(),
run.gate.warnings.length().to_string(),
])
}
to_csv(grid)
}
///|
pub fn pipeline_gate_summary(gate : QualityGateResult) -> String {
if gate.passed {
"quality gate passed: score=\{gate.score.to_string()}"
} else {
"quality gate failed: \{gate.reasons.length().to_string()} blocking reasons"
}
}
///|
pub fn pipeline_gate_feature_vector(gate : QualityGateResult) -> Array[Double] {
[
gate.score,
gate.clean_ratio,
gate.artifact_ratio,
gate.gap_ratio,
pipeline_bool_double(gate.passed),
gate.reasons.length().to_double(),
gate.warnings.length().to_double(),
]
}
///|
pub fn pipeline_stage_csv(run : HrvPipelineRun) -> String {
let grid = [
[
"ordinal", "stage", "status", "input_count", "output_count", "quality", "duration_ms",
"message",
],
]
for stage in run.stages {
grid.push([
stage.ordinal.to_string(),
pipeline_stage_name(stage.kind),
pipeline_status_name(stage.status),
stage.input_count.to_string(),
stage.output_count.to_string(),
stage.quality.to_string(),
stage.duration_ms.to_string(),
stage.message,
])
}
to_csv(grid)
}
///|
pub fn pipeline_policy_with_minimum_samples(
policy : QualityGatePolicy,
samples : Int,
) -> QualityGatePolicy {
{
minimum_samples: samples.max(1),
minimum_clean_ratio: policy.minimum_clean_ratio,
minimum_signal_quality: policy.minimum_signal_quality,
maximum_artifact_ratio: policy.maximum_artifact_ratio,
maximum_gap_ratio: policy.maximum_gap_ratio,
require_frequency: policy.require_frequency,
require_nonlinear: policy.require_nonlinear,
}
}
///|
pub fn pipeline_policy_is_conservative(policy : QualityGatePolicy) -> Bool {
policy.minimum_samples >= 32 &&
policy.minimum_clean_ratio >= 0.70 &&
policy.maximum_artifact_ratio <= 0.30
}