///|
/// A deterministic pipeline stage with explicit input and output metadata.
pub struct DataStage {
name : String
input_count : Int
output_count : Int
dropped_count : Int
score : Double
notes : Array[String]
}
///|
pub struct DataPipeline {
name : String
stages : Array[DataStage]
mut current : Array[Double]
mut completed : Bool
}
///|
pub struct DataPipelineResult {
data : Array[Double]
stages : Array[DataStage]
score : Double
completed : Bool
}
///|
pub fn data_stage(
name : String,
input_count : Int,
output_count : Int,
dropped_count : Int,
score : Double,
notes : Array[String],
) -> DataStage {
{ name, input_count, output_count, dropped_count, score, notes }
}
///|
pub fn data_pipeline(name : String, data : Array[Double]) -> DataPipeline {
{ name, stages: [], current: data.copy(), completed: false }
}
///|
fn pipeline_append(pipeline : DataPipeline, stage : DataStage) -> Unit {
pipeline.stages.push(stage)
}
///|
pub fn pipeline_validate(pipeline : DataPipeline) -> DataStage {
let report = validation_summary(pipeline.current, pipeline.name)
let score = report.score
let notes = validation_report_lines(report)
let stage = data_stage(
"validate",
pipeline.current.length(),
pipeline.current.length(),
0,
score,
notes,
)
pipeline_append(pipeline, stage)
stage
}
///|
pub fn pipeline_impute_median(pipeline : DataPipeline) -> DataStage {
let before = pipeline.current.length()
let missing = imputation_missing_count(pipeline.current)
pipeline.current = imputation_median(pipeline.current)
let stage = data_stage(
"impute-median",
before,
pipeline.current.length(),
missing,
if before == 0 {
1.0
} else {
1.0 - missing.to_double() / before.to_double()
},
[],
)
pipeline_append(pipeline, stage)
stage
}
///|
pub fn pipeline_clip(
pipeline : DataPipeline,
lower : Double,
upper : Double,
) -> DataStage {
let before = pipeline.current.length()
pipeline.current = transform_clip_array(pipeline.current, lower, upper)
let stage = data_stage("clip", before, pipeline.current.length(), 0, 1.0, [
"lower=" + lower.to_string(),
"upper=" + upper.to_string(),
])
pipeline_append(pipeline, stage)
stage
}
///|
pub fn pipeline_denoise(
pipeline : DataPipeline,
rule : SignalRule,
) -> DataStage {
let before = pipeline.current.length()
pipeline.current = signal_denoise(pipeline.current, rule)
let score = if before == 0 {
1.0
} else {
1.0 - transform_clip(signal_noise_estimate(pipeline.current), 0.0, 1.0)
}
let stage = data_stage(
"denoise",
before,
pipeline.current.length(),
0,
score,
[],
)
pipeline_append(pipeline, stage)
stage
}
///|
pub fn pipeline_standardize(pipeline : DataPipeline) -> DataStage {
let before = pipeline.current.length()
pipeline.current = transform_robust_scale(pipeline.current)
let stage = data_stage(
"robust-standardize",
before,
pipeline.current.length(),
0,
1.0,
[],
)
pipeline_append(pipeline, stage)
stage
}
///|
pub fn pipeline_downsample(
pipeline : DataPipeline,
target_size : Int,
) -> DataStage {
let before = pipeline.current.length()
pipeline.current = aggregate_downsample_median(pipeline.current, target_size)
let dropped = if before > pipeline.current.length() {
before - pipeline.current.length()
} else {
0
}
let score = if before == 0 {
1.0
} else {
pipeline.current.length().to_double() / before.to_double()
}
let stage = data_stage(
"downsample",
before,
pipeline.current.length(),
dropped,
score,
[],
)
pipeline_append(pipeline, stage)
stage
}
///|
pub fn pipeline_add_forecast_features(
pipeline : DataPipeline,
config : FeatureConfig,
) -> FeatureFrame {
feature_frame(pipeline.current, config)
}
///|
pub fn pipeline_finish(pipeline : DataPipeline) -> DataPipelineResult {
pipeline.completed = true
let mut total = 0.0
for stage in pipeline.stages {
total += stage.score
}
let score = if pipeline.stages.length() == 0 {
1.0
} else {
total / pipeline.stages.length().to_double()
}
{
data: pipeline.current.copy(),
stages: pipeline.stages.copy(),
score,
completed: pipeline.completed,
}
}
///|
pub fn pipeline_run_basic(
name : String,
data : Array[Double],
lower : Double,
upper : Double,
) -> DataPipelineResult {
let pipeline = data_pipeline(name, data)
let _ = pipeline_validate(pipeline)
let _ = pipeline_clip(pipeline, lower, upper)
pipeline_finish(pipeline)
}
///|
pub fn pipeline_run_robust(
name : String,
data : Array[Double],
lower : Double,
upper : Double,
window : Int,
) -> DataPipelineResult {
let pipeline = data_pipeline(name, data)
let _ = pipeline_validate(pipeline)
let _ = pipeline_impute_median(pipeline)
let _ = pipeline_clip(pipeline, lower, upper)
let _ = pipeline_denoise(pipeline, signal_rule(window, 3.5, 0.5, 2))
pipeline_finish(pipeline)
}
///|
pub fn pipeline_stage_names(result : DataPipelineResult) -> Array[String] {
let names = []
for stage in result.stages {
names.push(stage.name)
}
names
}
///|
pub fn pipeline_stage_scores(result : DataPipelineResult) -> Array[Double] {
let scores = []
for stage in result.stages {
scores.push(stage.score)
}
scores
}
///|
pub fn pipeline_stage_rows(result : DataPipelineResult) -> Array[Array[String]] {
let rows = [["stage", "input", "output", "dropped", "score"]]
for stage in result.stages {
rows.push([
stage.name,
stage.input_count.to_string(),
stage.output_count.to_string(),
stage.dropped_count.to_string(),
stage.score.to_string(),
])
}
rows
}
///|
pub fn pipeline_stage_table(result : DataPipelineResult) -> String {
let lines = []
for row in pipeline_stage_rows(result) {
lines.push(row.join("\t"))
}
lines.join("\n")
}
///|
pub fn pipeline_result_lines(result : DataPipelineResult) -> Array[String] {
[
"completed=" + result.completed.to_string(),
"count=" + result.data.length().to_string(),
"stages=" + result.stages.length().to_string(),
"score=" + result.score.to_string(),
]
}
///|
pub fn pipeline_result_string(result : DataPipelineResult) -> String {
pipeline_result_lines(result).join("\n")
}
///|
pub fn pipeline_quality_gate(
result : DataPipelineResult,
threshold : Double,
) -> Bool {
result.completed && result.score >= threshold && result.data.length() > 0
}
///|
pub fn pipeline_dropped_fraction(result : DataPipelineResult) -> Double {
let mut dropped = 0
let mut input = 0
for stage in result.stages {
dropped += stage.dropped_count
input += stage.input_count
}
if input == 0 {
0.0
} else {
dropped.to_double() / input.to_double()
}
}
///|
pub fn pipeline_reproducible(name : String, data : Array[Double]) -> Bool {
let left = pipeline_run_robust(name, data, -1.0e6, 1.0e6, 5)
let right = pipeline_run_robust(name, data, -1.0e6, 1.0e6, 5)
left.data == right.data &&
left.score == right.score &&
pipeline_stage_names(left) == pipeline_stage_names(right)
}
///|
pub fn data_pipeline_compare(
left : DataPipelineResult,
right : DataPipelineResult,
) -> Array[Double] {
[
left.score,
right.score,
right.score - left.score,
left.data.length().to_double(),
right.data.length().to_double(),
]
}
///|
pub fn pipeline_apply_transform(
result : DataPipelineResult,
rule : TransformRule,
) -> DataPipelineResult {
{ ..result, data: transform_apply_rule(result.data, rule) }
}
///|
pub fn pipeline_apply_signal(
result : DataPipelineResult,
rule : SignalRule,
) -> DataPipelineResult {
{ ..result, data: signal_denoise(result.data, rule) }
}
///|
pub fn pipeline_summary(result : DataPipelineResult) -> Array[Double] {
[
result.score,
pipeline_dropped_fraction(result),
result.data.length().to_double(),
result.stages.length().to_double(),
if result.completed {
1.0
} else {
0.0
},
]
}