#![allow(clippy::needless_range_loop)]
use std::env;
use std::fs;
use std::process;
use struktura::{analyze, bootstrap_alpha, health_check, prove_structure, HealthVerdict, LawQuality, StructuralLaw};
const NORMAL_SAMPLES: &str = include_str!("../../data/normal_sample.csv");
const FAULT_SAMPLES: &str = include_str!("../../data/fault_sample.csv");
fn read_input(path: &str) -> Vec<f64> {
if path == "-" {
return read_stdin();
}
if !std::path::Path::new(path).exists() {
eprintln!("Error: file not found: {}", path);
process::exit(1);
}
let content = fs::read_to_string(path).unwrap_or_else(|e| {
eprintln!("Error reading {}: {}", path, e);
process::exit(1);
});
parse_values_col(&content, resolve_col(&content))
}
fn read_csv(path: &str) -> Vec<f64> { read_input(path) }
fn read_stdin() -> Vec<f64> {
use std::io::Read;
let mut buf = String::new();
std::io::stdin().read_to_string(&mut buf).unwrap_or_else(|e| {
eprintln!("Error reading stdin: {}", e);
process::exit(1);
});
parse_values_col(&buf, resolve_col(&buf))
}
fn resolve_col(content: &str) -> Option<usize> {
let args: Vec<String> = env::args().collect();
let spec = args.iter().position(|a| a == "--col").and_then(|i| args.get(i + 1))?;
if let Ok(n) = spec.parse::<usize>() {
return Some(n);
}
let header = content.lines().find(|l| !l.is_empty() && !l.starts_with('#'))?;
let idx = header.split(',').position(|h| h.trim().eq_ignore_ascii_case(spec));
if idx.is_none() {
eprintln!("Error: column '{}' not found in header: {}", spec, header);
process::exit(1);
}
idx
}
fn read_text_input(path: &str) -> String {
if path == "-" {
use std::io::Read;
let mut buf = String::new();
std::io::stdin().read_to_string(&mut buf).unwrap_or_else(|e| {
eprintln!("Error reading stdin: {}", e);
process::exit(1);
});
return buf;
}
std::fs::read_to_string(path).unwrap_or_else(|e| {
eprintln!("Error reading {}: {}", path, e);
process::exit(1);
})
}
fn parse_values(content: &str) -> Vec<f64> {
parse_values_col(content, None)
}
fn parse_values_col(content: &str, column: Option<usize>) -> Vec<f64> {
content
.lines()
.filter(|l| !l.is_empty() && !l.starts_with('#'))
.filter_map(|l| {
let parts: Vec<&str> = l.split(',').collect();
let s = if let Some(col) = column {
parts.get(col).copied().unwrap_or("")
} else if parts.len() > 1 {
parts.last().copied().unwrap_or("")
} else {
l
};
s.trim().parse::<f64>().ok()
})
.collect()
}
struct ParsedCsv {
rows: Vec<Vec<f64>>,
col_names: Vec<String>,
ncols: usize,
}
fn parse_multi_csv(content: &str) -> ParsedCsv {
let delim = if content.lines().next().unwrap_or("").contains('\t') { '\t' }
else if content.lines().next().unwrap_or("").contains(';') { ';' }
else { ',' };
let mut rows: Vec<Vec<f64>> = Vec::new();
let mut col_mask: Option<Vec<bool>> = None;
let mut col_names: Vec<String> = Vec::new();
let mut header_seen = false;
for line in content.lines() {
let line = line.trim().trim_start_matches('\u{feff}');
if line.is_empty() || line.starts_with('#') { continue; }
let fields: Vec<&str> = line.split(delim).map(|s| s.trim().trim_matches('"')).collect();
if col_mask.is_none() {
let mask: Vec<bool> = fields.iter().map(|s| s.parse::<f64>().is_ok()).collect();
if mask.iter().all(|&m| !m) {
col_names = fields.iter().map(|s| s.to_string()).collect();
header_seen = true;
continue;
}
if header_seen && !col_names.is_empty() {
col_names = col_names.iter().zip(mask.iter())
.filter(|(_, &m)| m).map(|(n, _)| n.clone()).collect();
}
col_mask = Some(mask);
}
let mask = col_mask.as_ref().unwrap();
let vals: Vec<f64> = fields.iter().zip(mask.iter())
.filter(|(_, &m)| m)
.filter_map(|(s, _)| s.parse().ok())
.collect();
if !vals.is_empty() { rows.push(vals); }
}
let ncols = rows.first().map(|r| r.len()).unwrap_or(0);
ParsedCsv { rows, col_names, ncols }
}
impl ParsedCsv {
#[allow(dead_code)]
fn ch_name(&self, idx: usize) -> String {
self.col_names.get(idx).cloned().unwrap_or_else(|| format!("ch{}", idx))
}
fn sensor_columns(&self) -> (Vec<usize>, Vec<String>) {
let n = self.rows.len();
let mut keep = Vec::new();
let mut dropped = Vec::new();
for ch in 0..self.ncols {
let col: Vec<f64> = self.rows.iter().map(|r| r.get(ch).copied().unwrap_or(f64::NAN)).collect();
let step = if n >= 3 { col[1] - col[0] } else { 0.0 };
let magnitude = col.iter().fold(0.0f64, |m, v| m.max(v.abs()));
let tol = 1e-6 * step.abs() + 8.0 * f64::EPSILON * magnitude;
let evenly_rising = n >= 3
&& step > 0.0
&& col.windows(2).all(|w| ((w[1] - w[0]) - step).abs() <= tol);
if evenly_rising && self.ncols > 1 {
dropped.push(self.ch_name(ch));
} else {
keep.push(ch);
}
}
(keep, dropped)
}
}
fn index_columns_note(dropped: &[String]) {
if !dropped.is_empty() {
eprintln!(" note: not monitoring {} (rises in even steps, looks like an index or timestamp)", dropped.join(", "));
}
}
fn quality_str(q: LawQuality) -> &'static str {
match q {
LawQuality::Exact => "EXACT",
LawQuality::Strong => "STRONG",
LawQuality::Good => "GOOD",
LawQuality::Approx => "APPROX",
LawQuality::Abstain => "ABSTAIN",
LawQuality::Insufficient => "INSUFFICIENT",
_ => "UNKNOWN",
}
}
fn verdict_color(v: HealthVerdict) -> (&'static str, &'static str) {
match v {
HealthVerdict::Healthy => ("\x1b[32m", "HEALTHY"),
HealthVerdict::Watch => ("\x1b[33m", "WATCH"),
HealthVerdict::Warning => ("\x1b[33;1m", "WARNING"),
HealthVerdict::Critical => ("\x1b[31;1m", "CRITICAL"),
_ => ("", "UNKNOWN"),
}
}
fn alpha_bar(alpha: f64, width: usize) -> String {
let clamped = alpha.clamp(0.0, 1.5);
let filled = ((clamped / 1.5) * width as f64) as usize;
let empty = width.saturating_sub(filled);
format!("{}{}", "#".repeat(filled), ".".repeat(empty))
}
fn print_visual(label: &str, law: &StructuralLaw, baseline: Option<f64>) {
let bar = alpha_bar(law.dfa.alpha, 30);
let q = quality_str(law.quality);
println!(" {} {:<20} alpha={:.3} R2={:.3} [{}]", bar, label, law.dfa.alpha, law.dfa.r_squared, q);
if let Some(b) = baseline {
let verdict = health_check(law, b);
let shift = law.dfa.alpha - b;
let (color, label) = verdict_color(verdict);
println!(" {:>30} shift={:+.3} \x1b[1m{}{}\x1b[0m", "", shift, color, label);
}
}
fn print_law_detail(law: &StructuralLaw) {
println!(" Samples: {}", law.n);
println!(" DFA alpha: {:.3} (R2={:.4})", law.dfa.alpha, law.dfa.r_squared);
println!(" Hurst H: {:.3}", law.hurst);
println!(" Kurtosis: {:.1}", law.kurtosis);
println!(" Mean: {:.4} Std: {:.4}", law.mean, law.std_dev);
println!(" Quality: {}", quality_str(law.quality));
}
fn main() {
let args: Vec<String> = env::args().collect();
if args.len() < 2 || args[1] == "--help" || args[1] == "-h" || args[1] == "help" {
println!();
println!(" \x1b[1mstruktura\x1b[0m - detect when a signal's structure changes");
println!();
println!(" USE IT ON YOUR DATA:");
println!(" struktura guard <file.csv> Monitor a CSV for anomalies (autonomous)");
println!(" struktura guard <file.csv> --baseline 500 Calibrate on first 500 rows");
println!(" struktura guard <file.csv> --json Machine-readable output");
println!(" cat stream.csv | struktura guard - Pipe from stdin");
println!(" struktura when <file.csv> Find WHEN something changed (changepoint detection)");
println!(" struktura when <file.csv> --truth 0-500,9000-9999 Score changepoints against known windows");
println!(" struktura check <file.csv> One-shot structural analysis");
println!(" struktura prove <file.csv> Bootstrap CI on alpha + shuffle proof of structure");
println!(" struktura scan <file_or_-> Auto-classify + trend + health");
println!();
println!(" DEMOS (zero setup, embedded data):");
println!(" struktura demo Bearing fault (CWRU data)");
println!(" struktura nasa SMAP satellite anomaly");
println!(" struktura voyager Voyager 1 magnetometer, 2021 vs 2022");
println!(" struktura mission Autonomous 24K-sample gauntlet");
println!();
println!(" DOMAIN:");
println!(" struktura text <file.txt> Writing rhythm");
println!(" struktura market <prices.csv> Financial regime detection");
println!(" struktura rhythm <timestamps.csv> Event timing (heartbeat/git/IoT)");
println!(" struktura compare <a.csv> <b.csv> Compare two signals");
println!(" struktura copilot-compare <file.csv> DFA vs boolean threshold (side-by-side)");
println!(" struktura bench Full benchmark with all fault types");
println!(" struktura benchmark-faults F1 scores across 6 telemetry fault types");
println!();
println!(" INPUT: CSV or one-value-per-line. Uses last column, or --col <index|name>.");
println!(" MORE: https://github.com/koscak-labs/struktura");
println!();
process::exit(0);
}
match args[1].as_str() {
"demo" => cmd_demo(),
"check" => cmd_check(&args),
"prove" => cmd_prove(&args),
"compare" => cmd_compare(&args),
"stamp" => cmd_stamp(&args),
"bench" => cmd_bench(),
"report" => cmd_report(&args),
"validate" => cmd_validate(&args),
"self-test" => cmd_self_test(),
"codegen" => cmd_codegen(&args),
"generate" => cmd_generate(&args),
"voyager" => cmd_voyager(),
"heliopause" => cmd_heliopause(),
"ims" => cmd_ims(),
"spacecraft" => cmd_spacecraft(),
"genome" => cmd_genome(&args),
"text" => cmd_text(&args),
"market" => cmd_market(&args),
"rhythm" => cmd_rhythm(&args),
"scan" => cmd_scan(&args),
"classify" | "what" => cmd_classify(&args),
"multifractal" | "mf" => cmd_multifractal(&args),
"health" => cmd_health(&args),
"batch" => cmd_batch(&args),
"fingerprint" | "fp" | "dna" => cmd_fingerprint(&args),
"watch" => cmd_watch(&args),
"oneline" => cmd_oneline(&args),
"alert" => cmd_alert(&args),
"benchmark-faults" | "bf" => cmd_benchmark_faults(),
"benchmark-telemetry" | "bt" => cmd_benchmark_telemetry(&args),
"monitor-perf" => cmd_monitor_perf(),
"monitor-real" => cmd_monitor_real(),
"generate-hybrid" => cmd_generate_hybrid(&args),
"mission" => cmd_mission(),
"redblue" => cmd_redblue(&args),
"evolve" => cmd_evolve(&args),
"evolve-real" => cmd_evolve_real(&args),
"smap" => cmd_smap(&args),
"nasa" => cmd_nasa(),
"rover" => cmd_rover(),
"guard" => cmd_guard(&args),
"copilot-compare" | "cc" => cmd_copilot_compare(&args),
"when" => cmd_when(&args),
"pipe" => cmd_pipe(&args),
"investigate" => cmd_investigate(&args),
"case" => cmd_case(&args),
"replay" => cmd_replay(&args),
"version" => println!("struktura {}", env!("CARGO_PKG_VERSION")),
other => {
eprintln!("Unknown command: {}", other);
eprintln!("Try: struktura demo");
process::exit(1);
}
}
}
fn cmd_demo() {
let normal = parse_values(NORMAL_SAMPLES);
let fault = parse_values(FAULT_SAMPLES);
let law_n = analyze(&normal);
let law_f = analyze(&fault);
let verdict = health_check(&law_f, law_n.dfa.alpha);
let shift = law_f.dfa.alpha - law_n.dfa.alpha;
let (color, verdict_label) = verdict_color(verdict);
println!();
println!(" \x1b[1mSTRUKTURA DEMO\x1b[0m");
println!(" Bearing fault detection from CWRU vibration data");
println!(" ================================================");
println!();
print_visual("Normal bearing", &law_n, None);
println!();
print_visual("Inner race FAULT", &law_f, Some(law_n.dfa.alpha));
println!();
println!(" ================================================");
println!(" Baseline alpha: {:.3} (normal bearing)", law_n.dfa.alpha);
println!(" Current alpha: {:.3} (faulted bearing)", law_f.dfa.alpha);
println!(" Structural shift: \x1b[1m{:+.3}\x1b[0m", shift);
println!();
println!(" Verdict: {}\x1b[1m{}\x1b[0m", color, verdict_label);
println!();
println!(" The bearing's vibration structure changed BEFORE");
println!(" any amplitude threshold would have fired.");
println!();
println!(" No training. No hyperparameters. Just math.");
println!(" https://crates.io/crates/struktura");
println!();
}
fn cmd_check(args: &[String]) {
if args.len() < 3 {
eprintln!("Usage: struktura check <file.csv> [file2.csv ...] [--baseline N]");
process::exit(1);
}
let files: Vec<&str> = args[2..].iter()
.map(|s| s.as_str())
.take_while(|s| !s.starts_with("--"))
.collect();
if files.len() > 1 {
cmd_check_multi(&files, args);
return;
}
let path = &args[2];
let data = if path == "-" { read_stdin() } else { read_csv(path) };
if data.len() < 20 {
eprintln!("Error: need >= 20 data points, got {}", data.len());
process::exit(1);
}
let law = analyze(&data);
let stamp_baseline = if path != "-" {
std::fs::read_to_string(path).ok().and_then(|c| read_stamp(&c))
} else { None };
let mut baseline: Option<f64> = stamp_baseline;
let mut threshold: Option<f64> = None;
let json_mode = args.iter().any(|a| a == "--json");
let quiet_mode = args.iter().any(|a| a == "--quiet");
let csv_mode = args.iter().any(|a| a == "--csv");
for i in 0..args.len() {
if args[i] == "--baseline" && i + 1 < args.len() {
baseline = args[i + 1].parse().ok();
}
if args[i] == "--threshold" && i + 1 < args.len() {
threshold = args[i + 1].parse().ok();
}
}
let _ = threshold;
if csv_mode {
let shift_s = baseline.map(|b| format!("{:.4}", law.dfa.alpha - b)).unwrap_or_default();
let verdict_s = baseline.map(|b| { let (_, l) = verdict_color(health_check(&law, b)); l.to_string() }).unwrap_or_default();
println!("{},{},{:.4},{:.4},{:.4},{},{},{}", path, law.n, law.dfa.alpha, law.dfa.r_squared, law.hurst, quality_str(law.quality), shift_s, verdict_s);
return;
}
if quiet_mode {
if let Some(b) = baseline {
let v = health_check(&law, b);
let (_, l) = verdict_color(v);
println!("{}", l);
} else {
println!("{}", quality_str(law.quality));
}
return;
}
if json_mode {
let verdict_s = baseline.map(|b| {
let v = health_check(&law, b);
let (_, l) = verdict_color(v);
l
});
let shift = baseline.map(|b| law.dfa.alpha - b);
println!("{{\"file\":\"{}\",\"n\":{},\"dfa_alpha\":{:.4},\"dfa_r2\":{:.4},\"hurst\":{:.4},\"kurtosis\":{:.2},\"quality\":\"{}\"{}{}}}",
path, law.n, law.dfa.alpha, law.dfa.r_squared, law.hurst, law.kurtosis, quality_str(law.quality),
shift.map(|s| format!(",\"shift\":{:.4}", s)).unwrap_or_default(),
verdict_s.map(|v| format!(",\"verdict\":\"{}\"", v)).unwrap_or_default()
);
return;
}
println!();
println!(" \x1b[1mSTRUKTURA\x1b[0m - structural health analysis");
println!(" File: {}", path);
println!(" ------------------------------------------------");
print_visual(path, &law, baseline);
println!();
print_law_detail(&law);
let alpha = law.dfa.alpha;
let interp = if alpha < 0.4 {
"anti-correlated (mean-reverting). successive changes tend to cancel."
} else if alpha < 0.55 {
"uncorrelated (white noise / random walk increments). no temporal pattern."
} else if alpha < 0.75 {
"weakly correlated. some persistence but decays fast."
} else if alpha < 1.05 {
"strongly correlated (1/f noise). long-range memory in the signal."
} else if alpha < 1.4 {
"highly persistent (fractional Brownian). trends persist across scales."
} else {
"unbounded persistence (non-stationary random walk or worse)."
};
println!();
println!(" interpretation: alpha = {:.3} -> {}", alpha, interp);
if law.kurtosis > 5.0 {
println!(" kurtosis = {:.1} -> heavy-tailed / bursty (expect spikes).", law.kurtosis);
}
if law.dfa.r_squared < 0.85 {
println!(" R2 = {:.3} -> scaling law is a weak fit. be cautious.", law.dfa.r_squared);
}
if let Some(b) = baseline {
let verdict = health_check(&law, b);
let (color, label) = verdict_color(verdict);
let shift = (law.dfa.alpha - b).abs();
let ci = bootstrap_alpha(&data, 20);
let se = (ci.ci_high - ci.ci_low) / (2.0 * 1.96);
let z = if se > 0.0 { shift / se } else { f64::INFINITY };
println!();
if data.len() >= 256 {
println!(" >>> {}{}\x1b[0m (z={:.1} against this signal's alpha spread ±{:.3})", color, label, z, se * 1.96);
} else {
println!(" >>> {}{}\x1b[0m", color, label);
}
}
if args.iter().any(|a| a == "--shuffle") {
let proof = prove_structure(&data);
println!(" ------------------------------------------------");
println!(" \x1b[1mShuffle proof:\x1b[0m");
println!(" Real alpha: {:.3} (R2={:.4})", proof.real_alpha, proof.real_r2);
println!(" Shuffled alpha: {:.3} (R2={:.4})", proof.shuffled_alpha, proof.shuffled_r2);
if proof.structure_confirmed {
println!(" Result: \x1b[32mCONFIRMED\x1b[0m, shuffling destroyed the structure");
} else {
println!(" Result: \x1b[33mINCONCLUSIVE\x1b[0m, signal may lack exploitable structure");
}
}
println!();
}
fn cmd_check_multi(files: &[&str], args: &[String]) {
let mut baseline: Option<f64> = None;
let mut threshold: Option<f64> = None;
for i in 0..args.len() {
if args[i] == "--baseline" && i + 1 < args.len() {
baseline = args[i + 1].parse().ok();
}
if args[i] == "--threshold" && i + 1 < args.len() {
threshold = args[i + 1].parse().ok();
}
}
let _ = threshold;
println!();
println!(" \x1b[1mSTRUKTURA\x1b[0m: batch analysis ({} files)", files.len());
println!(" ================================================");
println!();
println!(" | {:<30} | {:>6} | {:>7} | {:>6} | {:>8} | {:>8} |", "File", "N", "Alpha", "R2", "Quality", "Verdict");
println!(" |{:-<32}|{:-<8}|{:-<9}|{:-<8}|{:-<10}|{:-<10}|", "", "", "", "", "", "");
for &path in files {
let data = read_csv(path);
if data.len() < 20 {
println!(" | {:<30} | {:>6} | {:>7} | {:>6} | {:>8} | {:>8} |",
path, data.len(), "--", "--", "TOO FEW", "--");
continue;
}
let law = analyze(&data);
let verdict_s = baseline.map(|b| {
let (_, l) = verdict_color(health_check(&law, b));
l.to_string()
}).unwrap_or_else(|| "--".to_string());
let short_path: &str = path.rsplit('/').next().unwrap_or(path);
let short_path: &str = short_path.rsplit('\\').next().unwrap_or(short_path);
println!(" | {:<30} | {:>6} | {:>7.3} | {:>6.4} | {:>8} | {:>8} |",
short_path, law.n, law.dfa.alpha, law.dfa.r_squared, quality_str(law.quality), verdict_s);
}
println!();
}
fn cmd_compare(args: &[String]) {
if args.len() < 4 {
eprintln!("Usage: struktura compare <baseline.csv> <current.csv>");
process::exit(1);
}
let data_a = read_csv(&args[2]);
let data_b = read_csv(&args[3]);
let law_a = analyze(&data_a);
let law_b = analyze(&data_b);
println!();
println!(" \x1b[1mSTRUKTURA\x1b[0m - structural comparison");
println!(" ================================================");
print_visual(&args[2], &law_a, None);
println!();
print_visual(&args[3], &law_b, Some(law_a.dfa.alpha));
println!();
println!(" ------------------------------------------------");
println!(" BASELINE:");
print_law_detail(&law_a);
println!(" CURRENT:");
print_law_detail(&law_b);
let shift = law_b.dfa.alpha - law_a.dfa.alpha;
let verdict = health_check(&law_b, law_a.dfa.alpha);
let (color, label) = verdict_color(verdict);
let direction = if shift > 0.0 { "more correlated (more persistent)" }
else { "less correlated (more random)" };
println!();
println!(" shift: {:+.3}, the signal became {}", shift, direction);
if data_a.len() >= 256 && data_b.len() >= 256 {
let ca = bootstrap_alpha(&data_a, 20);
let cb = bootstrap_alpha(&data_b, 20);
let se_a = (ca.ci_high - ca.ci_low) / (2.0 * 1.96);
let se_b = (cb.ci_high - cb.ci_low) / (2.0 * 1.96);
let se = (se_a * se_a + se_b * se_b).sqrt();
let z = if se > 0.0 { shift.abs() / se } else { f64::INFINITY };
println!(" spread: baseline ±{:.3} current ±{:.3} (95%, quarter-length windows)", se_a * 1.96, se_b * 1.96);
if z >= 3.0 {
println!(" >>> {}{}\x1b[0m (z={:.1}: the shift is outside both signals' own variability)", color, label, z);
} else {
println!(" >>> {}{}\x1b[0m (z={:.1}: the shift is within the signals' own variability, treat as inconclusive)", color, label, z);
}
} else {
println!(" >>> {}{}\x1b[0m", color, label);
}
println!();
}
fn cmd_bench() {
let normal = parse_values(NORMAL_SAMPLES);
let fault = parse_values(FAULT_SAMPLES);
let law_n = analyze(&normal);
let law_f = analyze(&fault);
println!();
println!(" \x1b[1mSTRUKTURA BENCHMARK\x1b[0m");
println!(" ================================================");
println!();
println!(" | Condition | DFA alpha | R2 | Shift | Verdict |");
println!(" |--------------------|-----------|--------|--------|----------|");
let print_row = |name: &str, law: &StructuralLaw, baseline: Option<f64>| {
let (shift_str, verdict_str) = match baseline {
Some(b) => {
let s = law.dfa.alpha - b;
let v = health_check(law, b);
let (_, vl) = verdict_color(v);
(format!("{:+.3}", s), vl.to_string())
}
None => ("--".to_string(), "--".to_string()),
};
println!(
" | {:<18} | {:.3} | {:.4} | {:>6} | {:>8} |",
name, law.dfa.alpha, law.dfa.r_squared, shift_str, verdict_str
);
};
print_row("Normal", &law_n, None);
print_row("Inner race fault", &law_f, Some(law_n.dfa.alpha));
println!();
println!(" Data: CWRU Bearing Data Center (builtin sample)");
println!(" Algorithm: Detrended Fluctuation Analysis (Peng 1994)");
println!(" https://crates.io/crates/struktura");
println!();
}
fn cmd_report(args: &[String]) {
let files: Vec<&str> = args[2..].iter()
.map(|s| s.as_str())
.take_while(|s| !s.starts_with("--"))
.collect();
if files.is_empty() {
eprintln!("Usage: struktura report <file1.csv> [file2.csv ...]");
process::exit(1);
}
let mut baseline: Option<f64> = None;
for i in 0..args.len() {
if args[i] == "--baseline" && i + 1 < args.len() {
baseline = args[i + 1].parse().ok();
}
}
println!("# Struktura Report\n");
println!("Generated by struktura v{}\n", env!("CARGO_PKG_VERSION"));
println!("| File | N | DFA alpha | R2 | Quality | Verdict |");
println!("|------|---|-----------|-----|---------|---------|");
for &path in &files {
let data = read_csv(path);
if data.len() < 20 { println!("| {} | {} | -- | -- | TOO FEW | -- |", path, data.len()); continue; }
let law = analyze(&data);
let vs = baseline.map(|b| { let (_, l) = verdict_color(health_check(&law, b)); l.to_string() }).unwrap_or_else(|| "--".to_string());
let short: &str = path.rsplit('/').next().unwrap_or(path);
let short: &str = short.rsplit('\\').next().unwrap_or(short);
println!("| {} | {} | {:.3} | {:.4} | {} | {} |", short, law.n, law.dfa.alpha, law.dfa.r_squared, quality_str(law.quality), vs);
}
println!("\nAlgorithm: Detrended Fluctuation Analysis (Peng et al., 1994)");
}
fn cmd_validate(args: &[String]) {
if args.len() < 3 {
eprintln!("Usage: struktura validate <file.csv>");
process::exit(1);
}
let path = &args[2];
if !std::path::Path::new(path).exists() {
println!("FAIL: file not found: {}", path);
process::exit(1);
}
let data = read_csv(path);
if data.is_empty() {
println!("FAIL: no numeric values found in {}", path);
process::exit(1);
}
if data.len() < 20 {
println!("WARN: only {} values (need >= 20 for analysis, >= 64 for DFA)", data.len());
} else if data.len() < 64 {
println!("WARN: {} values (DFA needs >= 64, ACR works with >= 20)", data.len());
} else {
println!("OK: {} numeric values, ready for analysis", data.len());
}
}
fn cmd_self_test() {
println!(" STRUKTURA SELF-TEST v{}", env!("CARGO_PKG_VERSION"));
println!(" ========================================");
let mut pass = 0;
let mut fail = 0;
let normal = parse_values(NORMAL_SAMPLES);
let fault = parse_values(FAULT_SAMPLES);
let law_n = analyze(&normal);
let law_f = analyze(&fault);
let verdict = health_check(&law_f, law_n.dfa.alpha);
if verdict == HealthVerdict::Critical {
println!(" [PASS] Demo data: fault detected as CRITICAL");
pass += 1;
} else {
println!(" [FAIL] Demo data: expected CRITICAL, got {:?}", verdict);
fail += 1;
}
if law_n.dfa.r_squared > 0.9 && law_f.dfa.r_squared > 0.9 {
println!(" [PASS] Both signals R2 > 0.9");
pass += 1;
} else {
println!(" [FAIL] R2 too low: normal={:.3} fault={:.3}", law_n.dfa.r_squared, law_f.dfa.r_squared);
fail += 1;
}
let proof = prove_structure(&normal);
if proof.structure_confirmed {
println!(" [PASS] Shuffle proof: structure CONFIRMED (real={:.3} shuffled={:.3})", proof.real_alpha, proof.shuffled_alpha);
pass += 1;
} else {
println!(" [FAIL] Shuffle proof: INCONCLUSIVE");
fail += 1;
}
let constant: Vec<f64> = vec![5.0; 100];
let law_c = analyze(&constant);
if law_c.quality == LawQuality::Abstain {
println!(" [PASS] Constant signal: Abstain (no panic)");
pass += 1;
} else {
println!(" [FAIL] Constant signal: expected Abstain, got {:?}", law_c.quality);
fail += 1;
}
let mut with_nan = normal.clone();
with_nan[0] = f64::NAN;
with_nan[10] = f64::INFINITY;
let law_nan = analyze(&with_nan);
if law_nan.n > 0 && law_nan.quality != LawQuality::Insufficient {
println!(" [PASS] NaN/Inf filtered: {} samples analyzed", law_nan.n);
pass += 1;
} else {
println!(" [FAIL] NaN handling failed");
fail += 1;
}
println!(" ========================================");
println!(" {}/{} passed", pass, pass + fail);
if fail > 0 {
process::exit(1);
}
}
fn cmd_codegen(args: &[String]) {
let mut window: usize = 256;
let mut threshold: f64 = 0.08;
for i in 0..args.len() {
if args[i] == "--window" && i + 1 < args.len() {
window = args[i + 1].parse().unwrap_or(256);
}
if args[i] == "--threshold" && i + 1 < args.len() {
threshold = args[i + 1].parse().unwrap_or(0.08);
}
}
let name = args.iter().position(|a| a == "--name")
.and_then(|i| args.get(i + 1))
.map(|s| s.as_str())
.unwrap_or("DfaHealthMonitor");
if args.iter().any(|a| a == "--fprime") {
print!("{}", struktura::codegen::generate_fprime_component(name, window));
} else if args.iter().any(|a| a == "--cfs") {
print!("{}", struktura::codegen::generate_cfs_app(name, window));
} else {
print!("{}", struktura::codegen::generate_c_monitor(window, threshold));
}
}
fn cmd_voyager() {
use struktura::space::voyager_demo;
let result = voyager_demo();
println!();
println!(" \x1b[1mVOYAGER 1 STRUCTURAL HEALTH ANALYSIS\x1b[0m");
println!(" DFA scaling analysis of magnetometer data across");
println!(" 2021 vs 2022 (embedded data, works after cargo install)");
println!(" ====================================================================");
println!();
println!(" Data: NASA SPDF, L.F. Burlaga, VIM 48-second averages");
println!(" https://spdf.gsfc.nasa.gov/pub/data/voyager/voyager1/");
println!();
let bar_h = make_bar(result.healthy_alpha);
let bar_a = make_bar(result.anomaly_alpha);
println!(" {} 2021 (healthy) alpha={:.3} R²={:.4}", bar_h, result.healthy_alpha, result.healthy_r2);
println!(" {} 2022 (May-Jul) alpha={:.3} R²={:.4}", bar_a, result.anomaly_alpha, result.anomaly_r2);
let v_str = match result.verdict {
HealthVerdict::Healthy => "\x1b[32mHEALTHY\x1b[0m",
HealthVerdict::Watch => "\x1b[33mWATCH\x1b[0m",
HealthVerdict::Warning => "\x1b[31mWARNING\x1b[0m",
HealthVerdict::Critical => "\x1b[31;1mCRITICAL\x1b[0m",
_ => "UNKNOWN",
};
println!(" {:30} shift={:+.3} {}", "", result.shift, v_str);
println!();
println!(" ====================================================================");
println!(" Baseline alpha: {:.3} (2021 healthy operation)", result.healthy_alpha);
println!(" 2022 alpha: {:.3} (May-Jul 2022)", result.anomaly_alpha);
println!();
println!(" Year-over-year comparison only. Pre-anomaly vs during-anomaly 2022");
println!(" shows no significant shift (p=0.52), and these slices give z=1.5:");
println!(" inconclusive. This is not a detection of the AACS anomaly.");
println!(" (the CRITICAL label is a fixed alpha threshold, not significance)");
println!();
}
fn cmd_heliopause() {
use struktura::space::heliopause_demo;
let result = heliopause_demo();
println!();
println!(" \x1b[1mVOYAGER 1 HELIOPAUSE CROSSING\x1b[0m");
println!(" DFA structural analysis across the boundary of our solar system");
println!(" August 25, 2012, first human-made object to enter interstellar space");
println!(" ====================================================================");
println!();
println!(" Data: NASA SPDF, L.F. Burlaga, VIM 48-second magnetometer averages");
println!(" https://spdf.gsfc.nasa.gov/pub/data/voyager/voyager1/");
println!();
let bar_h = make_bar(result.helio_alpha);
let bar_i = make_bar(result.interstellar_alpha);
println!(" {} Heliosphere (days 1-200) alpha={:.3} R²={:.4}", bar_h, result.helio_alpha, result.helio_r2);
println!(" {} Interstellar (days 260-331) alpha={:.3} R²={:.4}", bar_i, result.interstellar_alpha, result.interstellar_r2);
let v_str = match result.verdict {
HealthVerdict::Healthy => "\x1b[32mHEALTHY\x1b[0m",
HealthVerdict::Watch => "\x1b[33mWATCH\x1b[0m",
HealthVerdict::Warning => "\x1b[31mWARNING\x1b[0m",
HealthVerdict::Critical => "\x1b[31;1mCRITICAL\x1b[0m",
_ => "UNKNOWN",
};
println!(" {:30} shift={:+.3} {}", "", result.shift, v_str);
println!();
println!(" ====================================================================");
println!(" Heliosphere: alpha={:.3} (sun's magnetic field, strong persistence)", result.helio_alpha);
println!(" Interstellar: alpha={:.3} (after the crossing)", result.interstellar_alpha);
println!();
println!(" Alpha differs across the crossing, but the subsampling z on the");
println!(" bundled slices is 0.6: inconclusive. A longer pre/post window is");
println!(" the open test.");
println!();
}
fn cmd_batch(args: &[String]) {
if args.len() < 3 {
eprintln!("Usage: struktura batch <file1.csv> <file2.csv> ... [--baseline <file>] [--json]");
eprintln!(" Analyze multiple signals at once. Perfect for CI/CD monitoring.");
process::exit(1);
}
use struktura::classify::classify;
let mut files = Vec::new();
let mut baseline_file: Option<String> = None;
let mut json_mode = false;
let mut i = 2;
while i < args.len() {
match args[i].as_str() {
"--baseline" if i + 1 < args.len() => { baseline_file = Some(args[i + 1].clone()); i += 2; }
"--json" => { json_mode = true; i += 1; }
f => { files.push(f.to_string()); i += 1; }
}
}
let baseline = baseline_file.as_ref().map(|f| read_input(f));
let baseline_alpha = baseline.as_ref().map(|b| analyze(b).dfa.alpha);
if json_mode {
print!("[");
for (idx, f) in files.iter().enumerate() {
let values = read_input(f);
let law = analyze(&values);
let cls = classify(&values);
print!("{{\"file\":\"{}\",\"alpha\":{:.4},\"r_squared\":{:.4},\"quality\":\"{}\",\"type\":\"{}\",\"n\":{}",
f.replace('\\', "/"), law.dfa.alpha, law.dfa.r_squared, law.quality, cls.signal_type, law.n);
if let Some(ba) = baseline_alpha {
let shift = law.dfa.alpha - ba;
let verdict = health_check(&law, ba);
print!(",\"shift\":{:.4},\"verdict\":\"{}\"", shift, verdict);
}
print!("}}");
if idx < files.len() - 1 { print!(","); }
}
println!("]");
return;
}
println!();
println!(" \x1b[1mBATCH ANALYSIS\x1b[0m");
println!(" ====================================================================");
println!();
println!(" {:<40} {:>7} {:>7} {:>7} TYPE", "FILE", "α", "R²", "n");
println!(" {}", "─".repeat(80));
let mut any_alert = false;
for f in &files {
let values = read_input(f);
let law = analyze(&values);
let cls = classify(&values);
let short_name: String = f.chars().rev().take(38).collect::<String>().chars().rev().collect();
if let Some(ba) = baseline_alpha {
let verdict = health_check(&law, ba);
let (color, label) = verdict_color(verdict);
if verdict != HealthVerdict::Healthy { any_alert = true; }
println!(" {:<40} {:>7.3} {:>7.3} {:>7} {} {}{}\x1b[0m",
short_name, law.dfa.alpha, law.dfa.r_squared, law.n, cls.signal_type, color, label);
} else {
println!(" {:<40} {:>7.3} {:>7.3} {:>7} {}",
short_name, law.dfa.alpha, law.dfa.r_squared, law.n, cls.signal_type);
}
}
println!();
if any_alert {
println!(" \x1b[31;1m⚠ STRUCTURAL CHANGES DETECTED, see verdicts above\x1b[0m");
} else if baseline_alpha.is_some() {
println!(" \x1b[32m✓ All signals within baseline tolerance\x1b[0m");
}
println!();
}
fn cmd_alert(args: &[String]) {
if args.len() < 4 {
eprintln!("Usage: struktura alert <current> --baseline <healthy> [--threshold <0.08>]");
eprintln!(" Exit code 0 = healthy, 1 = degraded. For cron/systemd/CI.");
eprintln!(" Example: struktura alert now.csv --baseline ref.csv || send-slack-alert");
process::exit(2);
}
let mut baseline_file = String::new();
let mut threshold = 0.08f64;
let mut i = 3;
while i < args.len() {
match args[i].as_str() {
"--baseline" if i + 1 < args.len() => { baseline_file = args[i + 1].clone(); i += 2; }
"--threshold" if i + 1 < args.len() => { threshold = args[i + 1].parse().unwrap_or(0.08); i += 2; }
_ => { i += 1; }
}
}
if baseline_file.is_empty() {
eprintln!("Error: --baseline is required");
process::exit(2);
}
let current = read_input(&args[2]);
let baseline = read_input(&baseline_file);
let result = struktura::compare(&baseline, ¤t);
let degraded = result.shift.abs() > threshold;
eprintln!("struktura: α={:.3} shift={:+.3} {} (threshold={:.3})",
result.current_alpha, result.shift, result.verdict, threshold);
process::exit(if degraded { 1 } else { 0 });
}
fn cmd_oneline(args: &[String]) {
if args.len() < 3 {
eprintln!("Usage: struktura oneline <file_or_-> [--baseline <file>]");
eprintln!(" One-line summary for logs, alerts, Slack webhooks.");
process::exit(1);
}
use struktura::classify::classify;
let values = read_input(&args[2]);
let law = analyze(&values);
let cls = classify(&values);
let mut baseline_file: Option<String> = None;
let mut i = 3;
while i < args.len() {
if args[i] == "--baseline" && i + 1 < args.len() {
baseline_file = Some(args[i + 1].clone());
i += 2;
} else { i += 1; }
}
if let Some(ref bf) = baseline_file {
let baseline = read_input(bf);
let result = struktura::compare(&baseline, &values);
println!("struktura: α={:.3} R²={:.3} {} shift={:+.3} {} n={}",
law.dfa.alpha, law.dfa.r_squared, cls.signal_type, result.shift, result.verdict, law.n);
} else {
println!("struktura: α={:.3} R²={:.3} {} n={}",
law.dfa.alpha, law.dfa.r_squared, cls.signal_type, law.n);
}
}
fn cmd_watch(args: &[String]) {
if args.len() < 3 {
eprintln!("Usage: struktura watch <file.csv> [--interval <secs>] [--baseline <file>]");
eprintln!(" Re-analyzes a file every N seconds. Shows live structural changes.");
eprintln!(" Ctrl+C to stop.");
process::exit(1);
}
let file = &args[2];
let mut interval_secs = 5u64;
let mut baseline_file: Option<String> = None;
let mut i = 3;
while i < args.len() {
match args[i].as_str() {
"--interval" if i + 1 < args.len() => {
interval_secs = args[i + 1].parse().unwrap_or(5);
i += 2;
}
"--baseline" if i + 1 < args.len() => {
baseline_file = Some(args[i + 1].clone());
i += 2;
}
_ => { i += 1; }
}
}
let baseline_alpha = baseline_file.as_ref().map(|f| {
let data = read_input(f);
analyze(&data).dfa.alpha
});
println!();
println!(" \x1b[1mSTRUKTURA WATCH\x1b[0m: monitoring {} every {}s", file, interval_secs);
if let Some(ba) = baseline_alpha {
println!(" Baseline alpha: {:.3}", ba);
}
println!(" Press Ctrl+C to stop");
println!(" ────────────────────────────────────────────────────────────");
let mut prev_alpha = 0.0;
loop {
let values = read_input(file);
let law = analyze(&values);
let cls = struktura::classify::classify(&values);
let now = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_secs())
.unwrap_or(0);
let ts = format!("{:02}:{:02}:{:02}",
(now / 3600) % 24, (now / 60) % 60, now % 60);
let bar = alpha_bar(law.dfa.alpha, 20);
let delta = if prev_alpha != 0.0 {
format!("Δ={:+.4}", law.dfa.alpha - prev_alpha)
} else {
"Δ=---".to_string()
};
if let Some(ba) = baseline_alpha {
let verdict = health_check(&law, ba);
let (color, label) = verdict_color(verdict);
println!(" [{}] {} α={:.3} R²={:.3} {} {} {}{}\x1b[0m",
ts, bar, law.dfa.alpha, law.dfa.r_squared, delta, cls.signal_type, color, label);
} else {
println!(" [{}] {} α={:.3} R²={:.3} {} {}",
ts, bar, law.dfa.alpha, law.dfa.r_squared, delta, cls.signal_type);
}
prev_alpha = law.dfa.alpha;
std::thread::sleep(std::time::Duration::from_secs(interval_secs));
}
}
fn cmd_fingerprint(args: &[String]) {
if args.len() < 3 {
eprintln!("Usage: struktura dna <file1> [file2] ... (aliases: fingerprint, fp)");
eprintln!(" Structural DNA, the unique identity of each signal.");
process::exit(1);
}
use struktura::fingerprint::fingerprint;
println!();
println!(" \x1b[1mSTRUCTURAL DNA\x1b[0m");
println!(" ====================================================================");
println!();
let fps: Vec<_> = args[2..].iter().map(|f| {
let values = read_input(f);
let fp = fingerprint(&values);
(f.clone(), fp)
}).collect();
for (f, fp) in &fps {
println!(" {}", f);
println!(" {fp}");
println!();
}
if fps.len() >= 2 {
println!(" \x1b[1mSimilarity Matrix\x1b[0m");
for i in 0..fps.len() {
for j in (i+1)..fps.len() {
let d = fps[i].1.distance(&fps[j].1);
let sim = if d < 0.05 { "\x1b[32mNEAR-IDENTICAL\x1b[0m" }
else if d < 0.2 { "\x1b[33mSIMILAR\x1b[0m" }
else { "\x1b[36mDIFFERENT\x1b[0m" };
let short_i: String = fps[i].0.chars().rev().take(20).collect::<String>().chars().rev().collect();
let short_j: String = fps[j].0.chars().rev().take(20).collect::<String>().chars().rev().collect();
println!(" {} ↔ {} distance={:.3} {sim}", short_i, short_j, d);
}
}
println!();
}
}
fn cmd_classify(args: &[String]) {
if args.len() < 3 {
eprintln!("Usage: struktura classify <file_or_-> [--json] (alias: what)");
eprintln!(" What kind of signal is this? White noise? Brownian? 1/f?");
process::exit(1);
}
use struktura::classify::classify;
let json_mode = args.iter().any(|a| a == "--json");
let values = read_input(&args[2]);
let result = classify(&values);
if json_mode {
println!("{{\"alpha\":{:.4},\"r_squared\":{:.4},\"type\":\"{}\",\"description\":\"{}\",\"n\":{}}}",
result.alpha, result.r_squared, result.signal_type, result.description, result.law.n);
return;
}
println!();
println!(" \x1b[1mSIGNAL CLASSIFICATION\x1b[0m");
println!(" ====================================================================");
println!();
let bar = alpha_bar(result.alpha, 30);
println!(" {} α={:.3} R²={:.4}", bar, result.alpha, result.r_squared);
println!();
println!(" Type: \x1b[1m{}\x1b[0m", result.signal_type);
println!(" {}", result.description);
println!();
println!(" ┌─────────────────────────────────────────────────────┐");
println!(" │ α<0.3 0.5 0.6-0.85 0.85-1.15 1.15-1.65 │");
println!(" │ anti white correlated 1/f brownian │");
println!(" │ ─────────────────▲─────────────────────────────── │");
let pos = ((result.alpha.clamp(0.0, 2.0) / 2.0) * 50.0) as usize;
let indicator = format!(" │ {}▲", " ".repeat(pos.min(50)));
println!("{}{}│", indicator, " ".repeat(53usize.saturating_sub(indicator.len() - 5)));
println!(" └─────────────────────────────────────────────────────┘");
println!();
}
fn cmd_genome(args: &[String]) {
use struktura::genome::genome_structure;
if args.len() < 3 {
eprintln!("Usage: struktura genome <file.fa> [--window 1000] [--profile 4096]");
process::exit(1);
}
let file = &args[2];
let content = std::fs::read_to_string(file).unwrap_or_else(|e| {
eprintln!("Error reading {}: {}", file, e);
process::exit(1);
});
let window: usize = args.iter().position(|a| a == "--window")
.and_then(|i| args.get(i + 1))
.and_then(|s| s.parse().ok())
.unwrap_or(1000);
let profile_block: usize = args.iter().position(|a| a == "--profile")
.and_then(|i| args.get(i + 1))
.and_then(|s| s.parse().ok())
.unwrap_or(4096);
let result = genome_structure(&content, window, profile_block);
println!();
println!(" \x1b[1mGENOME STRUCTURAL ANALYSIS\x1b[0m");
println!(" DFA scaling analysis of GC-content across {}", file);
println!(" ====================================================================");
println!();
println!(" Bases: {} ({:.1} Mbp)", result.total_bases, result.total_bases as f64 / 1e6);
println!(" GC windows: {} ({}bp each)", result.gc_windows, result.window_size);
println!(" GC content: {:.1}%", result.gc_content * 100.0);
println!();
let bar = make_bar(result.overall_dfa.alpha);
println!(" {} overall α={:.3} R²={:.4}", bar, result.overall_dfa.alpha, result.overall_dfa.r_squared);
println!();
if !result.profile.is_empty() {
println!(" Structural profile ({:.0}kb blocks):", profile_block as f64 * window as f64 / 1000.0);
for r in &result.profile {
let b = make_bar(r.alpha);
println!(" {:.1}-{:.1}M {} α={:.3}", r.start_bp as f64/1e6, r.end_bp as f64/1e6, b, r.alpha);
}
println!();
}
println!(" ====================================================================");
println!(" α ≈ 1.0: strong 1/f organization (gene-rich, regulated)");
println!(" α ≈ 0.7: moderate (typical heterochromatin)");
println!(" α ≈ 0.5: near-random (structural desert)");
println!();
}
fn cmd_ims() {
use struktura::space::ims_demo;
let result = ims_demo();
println!();
println!(" \x1b[1mNASA IMS BEARING RUN-TO-FAILURE\x1b[0m");
println!(" DFA structural analysis across 984 recordings (7 days)");
println!(" Bearing 1 outer race fault, IMS/U. Cincinnati, NASA DASHlink");
println!(" ====================================================================");
println!();
println!(" Baseline (recordings 1-900): α≈{:.3} (stable, anti-correlated noise)", result.baseline_alpha);
println!(" First alarm (recording 970): α={:.3}", result.pre_failure_alpha);
println!(" Last recording (984): α={:.3}", result.failure_alpha);
println!();
println!(" Timeline:");
println!(" rec 1-900: α≈0.17 ▒▒▒▒▒▒▒▒▒▒▒▒▒▒▒▒▒▒▒▒▒▒▒▒▒▒▒ stable");
println!(" rec 970: α=0.30 ▒▒▒▒▒▒▒▒▒▒▒▒██");
println!(" rec 974: α=0.53 ▒▒▒▒▒▒▒▒▒▒▒▒██████████");
println!(" rec 981: α=0.47 ▒▒▒▒▒▒▒▒▒▒▒▒█████████");
println!(" rec 983: α=0.18 ▒▒▒▒▒▒");
println!(" rec 984: α=0.11 ▒▒▒");
println!();
println!(" ====================================================================");
println!(" DFA's alpha shifted (regime change) starting at recording 970, before");
println!(" the recording 984 failure. A plain RMS threshold on this same data");
println!(" trips earlier than that, so this is not an early-warning result,");
println!(" it demonstrates the same DFA core used in the heliopause demo.");
println!();
}
fn cmd_health(args: &[String]) {
if args.len() < 3 {
eprintln!("Usage: struktura health <current> --baseline <healthy>");
eprintln!(" Full diagnostic: DFA + MFDFA + statistics + verdict.");
process::exit(1);
}
use struktura::mfdfa::mfdfa;
let mut baseline_file: Option<String> = None;
let mut i = 3;
while i < args.len() {
if args[i] == "--baseline" && i + 1 < args.len() {
baseline_file = Some(args[i + 1].clone());
i += 2;
} else { i += 1; }
}
let current = read_input(&args[2]);
if current.len() < 64 {
eprintln!("Need at least 64 data points, got {}", current.len());
process::exit(1);
}
let law = analyze(¤t);
let qs = [-5.0, -3.0, -1.0, 0.0, 1.0, 2.0, 3.0, 5.0];
let spectrum = mfdfa(¤t, &qs);
let proof = struktura::prove_structure(¤t);
println!();
println!(" \x1b[1mFULL HEALTH DIAGNOSTIC\x1b[0m");
println!(" ====================================================================");
println!();
println!(" Signal: {} samples", current.len());
println!(" Mean: {:.4} Std: {:.4} Kurtosis: {:.2}", law.mean, law.std_dev, law.kurtosis);
println!();
println!(" \x1b[1mStructural Analysis (DFA)\x1b[0m");
let bar = make_bar(law.dfa.alpha);
println!(" {} α={:.3} R²={:.4} quality={}", bar, law.dfa.alpha, law.dfa.r_squared, law.quality);
println!(" Hurst exponent: {:.3}", law.hurst);
println!(" Structure proof: {} (real α={:.3} vs shuffled α={:.3})",
if proof.structure_confirmed { "\x1b[32mCONFIRMED\x1b[0m" } else { "\x1b[33mINCONCLUSIVE\x1b[0m" },
proof.real_alpha, proof.shuffled_alpha);
println!();
println!(" \x1b[1mMultifractal Spectrum (MFDFA)\x1b[0m");
println!(" h(2)={:.3} width={:.3} {}",
spectrum.h2, spectrum.width,
if spectrum.is_multifractal { "\x1b[35mMULTIFRACTAL\x1b[0m" } else { "monofractal" });
print!(" h(q): ");
for p in &spectrum.points { print!("{:.2} ", p.h_q); }
println!();
if current.len() >= 512 {
let trend = struktura::trend::alpha_trend(¤t, 256, 64);
let dir_str = match trend.direction {
struktura::trend::TrendDirection::Improving => "\x1b[32mIMPROVING\x1b[0m",
struktura::trend::TrendDirection::Stable => "\x1b[32mSTABLE\x1b[0m",
struktura::trend::TrendDirection::Degrading => "\x1b[31mDEGRADING\x1b[0m",
};
println!();
println!(" \x1b[1mTrend (is it getting worse?)\x1b[0m");
println!(" α {:.3} → {:.3} slope={:.6}/sample {dir_str}",
trend.alpha_start, trend.alpha_end, trend.slope);
}
if let Some(ref bf) = baseline_file {
let baseline = read_input(bf);
let result = struktura::compare(&baseline, ¤t);
let baseline_spectrum = mfdfa(&baseline, &qs);
let v_str = match result.verdict {
HealthVerdict::Healthy => "\x1b[32mHEALTHY\x1b[0m",
HealthVerdict::Watch => "\x1b[33mWATCH\x1b[0m",
HealthVerdict::Warning => "\x1b[31mWARNING\x1b[0m",
HealthVerdict::Critical => "\x1b[31;1mCRITICAL\x1b[0m",
_ => "UNKNOWN",
};
println!();
println!(" \x1b[1mvs Baseline\x1b[0m");
println!(" DFA shift: {:+.3} {v_str}", result.shift);
println!(" MFDFA width: {:.3} → {:.3} (Δ={:+.3})",
baseline_spectrum.width, spectrum.width, spectrum.width - baseline_spectrum.width);
}
println!();
println!(" ====================================================================");
println!();
}
fn cmd_multifractal(args: &[String]) {
if args.len() < 3 {
eprintln!("Usage: struktura multifractal <file_or_-> (alias: mf)");
eprintln!(" Reveals multi-scale structure, is the signal simple or complex?");
process::exit(1);
}
use struktura::mfdfa::mfdfa;
let values = read_input(&args[2]);
if values.len() < 64 {
eprintln!("Need at least 64 data points, got {}", values.len());
process::exit(1);
}
let qs = [-5.0, -3.0, -2.0, -1.0, 0.0, 1.0, 2.0, 3.0, 5.0];
let spectrum = mfdfa(&values, &qs);
println!();
println!(" \x1b[1mMULTIFRACTAL SPECTRUM\x1b[0m");
println!(" How does structure vary across scales?");
println!(" ====================================================================");
println!();
println!(" h(2) = {:.3} (equivalent to standard DFA alpha)", spectrum.h2);
println!(" Spectral width = {:.3}", spectrum.width);
println!();
println!(" q h(q) R²");
println!(" ─────────────────────────────────────────");
for p in &spectrum.points {
let bar_len = ((p.h_q * 40.0).round() as usize).min(40);
let bar = "#".repeat(bar_len);
let pad = ".".repeat(40 - bar_len);
println!(" {:+5.1} {bar}{pad} {:.3} {:.3}", p.q, p.h_q, p.r_squared);
}
println!();
let mf_str = if spectrum.is_multifractal {
"\x1b[35mMULTIFRACTAL\x1b[0m, different scales have different structure"
} else {
"\x1b[32mMONOFRACTAL\x1b[0m, uniform structure across all scales"
};
println!(" Verdict: {mf_str}");
println!(" Width {:.3} (wider = more complex multi-scale behavior)", spectrum.width);
println!();
}
fn cmd_scan(args: &[String]) {
if args.len() < 3 {
eprintln!("Usage: struktura scan <file_or_-> [--baseline <file>] [--json]");
eprintln!(" Auto-analyzes any signal. Pipe from stdin with -");
eprintln!(" Examples:");
eprintln!(" struktura scan sensor.csv");
eprintln!(" cat data.csv | struktura scan -");
eprintln!(" struktura scan current.csv --baseline healthy.csv");
eprintln!(" struktura scan data.csv --json # machine-readable output");
process::exit(1);
}
let values = read_input(&args[2]);
if values.is_empty() {
eprintln!("No numeric values found in input");
process::exit(1);
}
let mut baseline_file: Option<String> = None;
let mut json_mode = false;
let mut i = 3;
while i < args.len() {
match args[i].as_str() {
"--baseline" if i + 1 < args.len() => { baseline_file = Some(args[i + 1].clone()); i += 2; }
"--json" => { json_mode = true; i += 1; }
_ => { i += 1; }
}
}
use struktura::classify::classify;
let law = analyze(&values);
let cls = classify(&values);
if json_mode {
let baseline_info = baseline_file.as_ref().map(|bf| {
let baseline = read_input(bf);
struktura::compare(&baseline, &values)
});
print!("{{\"alpha\":{:.4},\"r_squared\":{:.4},\"quality\":\"{}\",\"type\":\"{}\",\"n\":{},\"mean\":{:.4},\"std\":{:.4},\"kurtosis\":{:.2},\"hurst\":{:.3}",
law.dfa.alpha, law.dfa.r_squared, law.quality, cls.signal_type, law.n, law.mean, law.std_dev, law.kurtosis, law.hurst);
if let Some(ref result) = baseline_info {
print!(",\"shift\":{:.4},\"verdict\":\"{}\"", result.shift, result.verdict);
}
println!("}}");
return;
}
println!();
println!(" \x1b[1mSTRUKTURA SCAN\x1b[0m");
println!(" ====================================================================");
println!();
let bar = alpha_bar(law.dfa.alpha, 30);
let q = quality_str(law.quality);
println!(" {} α={:.3} R²={:.4} [{q}] n={}", bar, law.dfa.alpha, law.dfa.r_squared, law.n);
println!(" {:30} type: \x1b[1m{}\x1b[0m", "", cls.signal_type);
println!(" {:30} mean={:.2} std={:.2} kurtosis={:.2} H={:.3}", "", law.mean, law.std_dev, law.kurtosis, law.hurst);
if values.len() >= 512 {
let trend = struktura::trend::alpha_trend(&values, 256, 64);
let dir_str = match trend.direction {
struktura::trend::TrendDirection::Improving => "\x1b[32mSTABLE/IMPROVING\x1b[0m",
struktura::trend::TrendDirection::Stable => "\x1b[32mSTABLE\x1b[0m",
struktura::trend::TrendDirection::Degrading => "\x1b[31mDEGRADING\x1b[0m",
};
println!(" {:30} trend: α {:.3}→{:.3} {dir_str}", "", trend.alpha_start, trend.alpha_end);
}
if let Some(ref bf) = baseline_file {
let baseline = read_input(bf);
let result = struktura::compare(&baseline, &values);
let (color, label) = verdict_color(result.verdict);
println!();
println!(" vs baseline ({bf}):");
println!(" {:30} shift={:+.3} {}{}\x1b[0m", "", result.shift, color, label);
}
println!();
println!(" ====================================================================");
println!();
}
fn cmd_market(args: &[String]) {
if args.len() < 3 {
eprintln!("Usage: struktura market <prices.csv>");
eprintln!(" CSV with one price per line (close prices). Need 200+ prices.");
process::exit(1);
}
use struktura::market::{returns_from_prices, regime_detect};
let prices = read_csv(&args[2]);
if prices.len() < 64 {
eprintln!("Need at least 64 prices, got {}", prices.len());
process::exit(1);
}
let returns = returns_from_prices(&prices);
let result = regime_detect(&returns);
println!();
println!(" \x1b[1mMARKET REGIME DETECTION\x1b[0m");
println!(" DFA on log-returns, what strategy works NOW");
println!(" ====================================================================");
println!();
let bar = make_bar(result.alpha);
let regime_str = match result.regime {
struktura::market::MarketRegime::Trending => "\x1b[32mTRENDING\x1b[0m (momentum works)",
struktura::market::MarketRegime::RandomWalk => "\x1b[33mRANDOM WALK\x1b[0m (no reliable edge)",
struktura::market::MarketRegime::MeanReverting => "\x1b[36mMEAN-REVERTING\x1b[0m (reversals work)",
};
println!(" {} α={:.3} R²={:.4} n={}", bar, result.alpha, result.r_squared, result.n_returns);
println!(" {:30} regime: {}", "", regime_str);
println!();
println!(" ====================================================================");
println!(" α > 0.6 = trending (momentum)");
println!(" α ≈ 0.5 = random walk (efficient market)");
println!(" α < 0.45 = mean-reverting (reversal strategies)");
println!();
}
fn cmd_rhythm(args: &[String]) {
if args.len() < 3 {
eprintln!("Usage: struktura rhythm <timestamps.csv>");
eprintln!(" One unix timestamp per line. Try: git log --format=%at | struktura rhythm -");
process::exit(1);
}
use struktura::rhythm::{intervals_from_timestamps, rhythm_analyze};
let timestamps = read_csv(&args[2]);
let intervals = intervals_from_timestamps(×tamps);
if intervals.len() < 64 {
eprintln!("Need at least 64 intervals (65 timestamps), got {}", intervals.len());
process::exit(1);
}
let result = rhythm_analyze(&intervals);
println!();
println!(" \x1b[1mEVENT RHYTHM ANALYSIS\x1b[0m");
println!(" DFA on inter-event intervals");
println!(" ====================================================================");
println!();
let bar = make_bar(result.alpha);
let rhythm_str = match result.rhythm {
struktura::rhythm::RhythmType::Bursty => "\x1b[35mBURSTY\x1b[0m (creative/flow state bursts)",
struktura::rhythm::RhythmType::Natural => "\x1b[32mNATURAL\x1b[0m (human work rhythm)",
struktura::rhythm::RhythmType::Random => "\x1b[33mRANDOM\x1b[0m (uncorrelated events)",
struktura::rhythm::RhythmType::Metronomic => "\x1b[36mMETRONOMIC\x1b[0m (automated/scheduled)",
};
println!(" {} α={:.3} R²={:.4} n={}", bar, result.alpha, result.r_squared, result.n_intervals);
println!(" {:30} rhythm: {}", "", rhythm_str);
println!(" {:30} mean interval: {:.1}s", "", result.mean_interval);
println!();
println!(" ====================================================================");
println!(" α > 0.7 = bursty (long correlated work sessions)");
println!(" α 0.55-0.7 = natural human rhythm");
println!(" α ≈ 0.5 = random/uncorrelated");
println!(" α < 0.45 = metronomic (bot/cron/automated)");
println!();
}
fn cmd_text(args: &[String]) {
if args.len() < 3 {
eprintln!("Usage: struktura text <file.txt> [file2.txt ...]");
process::exit(1);
}
use struktura::text::text_structure;
println!();
println!(" \x1b[1mTEXT STRUCTURE ANALYSIS\x1b[0m");
println!(" DFA on sentence-length sequences, measures writing rhythm");
println!(" ====================================================================");
println!();
println!(" Reference: human prose α≈0.7-0.8 | shuffled/mechanical α≈0.5");
println!();
for path in &args[2..] {
let text = read_text_input(path);
let result = text_structure(&text);
let bar = make_bar(result.dfa.alpha);
println!(" {} {}", bar, path);
println!(" {:30} sentences={} mean_len={:.0} α={:.3} R²={:.4}",
"", result.sentence_count, result.mean_sentence_len,
result.dfa.alpha, result.dfa.r_squared);
if result.sentence_count >= 64 && result.dfa.r_squared > 0.5 {
let tag = if result.dfa.alpha > 0.7 {
"\x1b[32mSTRONG RHYTHM\x1b[0m (human literary)"
} else if result.dfa.alpha > 0.55 {
"\x1b[33mMODERATE RHYTHM\x1b[0m"
} else {
"\x1b[36mUNIFORM/MECHANICAL\x1b[0m"
};
println!(" {:30} {}", "", tag);
} else if result.sentence_count < 64 {
println!(" {:30} \x1b[90m(need 200+ sentences for reliable DFA)\x1b[0m", "");
}
println!();
}
println!(" ====================================================================");
println!(" α > 0.6 = long-range correlations in sentence rhythm (human writing)");
println!(" α ≈ 0.5 = random/shuffled/uniform sentence lengths");
println!();
}
fn cmd_spacecraft() {
use struktura::space::spacecraft_demo;
println!();
println!(" \x1b[1mMULTI-CHANNEL SPACECRAFT HEALTH MONITOR\x1b[0m");
println!(" DFA structural analysis across 4 telemetry channels");
println!(" (RWA/BAT/THM synthetic; MAG = real Voyager 1 data)");
println!(" ====================================================================");
println!();
let results = spacecraft_demo();
for r in &results {
let bar = make_bar(r.current_alpha);
let v_str = match r.verdict {
HealthVerdict::Healthy => "\x1b[32mHEALTHY\x1b[0m",
HealthVerdict::Watch => "\x1b[33mWATCH\x1b[0m",
HealthVerdict::Warning => "\x1b[31mWARNING\x1b[0m",
HealthVerdict::Critical => "\x1b[31;1mCRITICAL\x1b[0m",
_ => "UNKNOWN",
};
println!(" {} [{}:{}]", bar, r.subsystem, r.channel_name);
println!(" {:30} alpha={:.3} baseline={:.3} shift={:+.3} R²={:.4} {}",
"", r.current_alpha, r.baseline_alpha, r.shift, r.r_squared, v_str);
println!();
}
println!(" ====================================================================");
println!(" Each channel: first half = baseline, second half = current period.");
println!(" DFA reads structure, not level: the mean can look normal while");
println!(" alpha shifts. Whether that shift is a fault needs a control run.");
println!();
}
fn make_bar(alpha: f64) -> String {
let filled = (alpha * 30.0).round() as usize;
let filled = if filled > 30 { 30 } else { filled };
let empty = 30 - filled;
format!("{}{}", "#".repeat(filled), ".".repeat(empty))
}
fn cmd_generate(args: &[String]) {
let mut db_path: Option<String> = None;
let mut target = "cfs";
let mut output_dir = "dfa_monitor_app".to_string();
let mut window: usize = 256;
let mut threshold: f64 = 0.08;
let mut i = 2;
while i < args.len() {
match args[i].as_str() {
"--db" if i + 1 < args.len() => { db_path = Some(args[i + 1].clone()); i += 2; }
"--cfs" => { target = "cfs"; i += 1; }
"--fprime" => { target = "fprime"; i += 1; }
"--rover" => { target = "rover"; i += 1; }
"--ros" | "--ros2" => { target = "ros"; i += 1; }
"--window" if i + 1 < args.len() => { window = args[i + 1].parse().unwrap_or(256); i += 2; }
"--threshold" if i + 1 < args.len() => { threshold = args[i + 1].parse().unwrap_or(0.08); i += 2; }
"-o" | "--output" if i + 1 < args.len() => { output_dir = args[i + 1].clone(); i += 2; }
"--help" => {
println!("struktura generate, generate complete cFS/F Prime DFA monitor apps");
println!();
println!(" Compatible with nasa/ogma db.json format. No Haskell required.");
println!();
println!(" OPTIONS:");
println!(" --db <db.json> Variable database (ogma-compatible)");
println!(" --cfs Generate cFS application (default)");
println!(" --fprime Generate F Prime component");
println!(" --ros Generate ROS 2 node");
println!(" --window <N> DFA window size (default: 256)");
println!(" --threshold <F> Shift threshold (default: 0.08)");
println!(" -o, --output <dir> Output directory (default: dfa_monitor_app)");
println!();
println!(" EXAMPLE:");
println!(" struktura generate --cfs --db channels.json -o my_monitor/");
println!();
println!(" The db.json format is compatible with nasa/ogma's variable database.");
println!(" See: https://github.com/nasa/ogma");
process::exit(0);
}
_ => { i += 1; }
}
}
let channels = if let Some(ref path) = db_path {
parse_db_json(path)
} else {
vec![Channel { name: "input_value".into(), c_type: "double".into(), topic: "SAMPLE_MID".into(), field: "payload".into(), msg_type: "sample_msg_t".into() }]
};
let out = std::path::Path::new(&output_dir);
fs::create_dir_all(out).unwrap_or_else(|e| {
eprintln!("Error creating {}: {}", output_dir, e);
process::exit(1);
});
match target {
"cfs" => generate_cfs_app_dir(out, &channels, window, threshold),
"fprime" => generate_fprime_dir(out, &channels, window, threshold),
"ros" => generate_ros_dir(out, &channels, window, threshold),
"rover" => {
let fpp = struktura::codegen::generate_fprime_rover();
let fpp_path = format!("{}/RoverHealth.fpp", output_dir);
std::fs::create_dir_all(&output_dir).ok();
std::fs::write(&fpp_path, &fpp).expect("write .fpp");
let c_wrapper = generate_rover_c_wrapper();
std::fs::write(format!("{}/RoverHealthImpl.c", output_dir), &c_wrapper).expect("write C wrapper");
println!("Generated rover F Prime component in {}/", output_dir);
println!(" RoverHealth.fpp , F Prime component definition");
println!(" RoverHealthImpl.c , C implementation (calls rover_flight staticlib)");
println!(" Link with: cargo build --release --lib --crate-type staticlib");
return;
}
_ => { eprintln!("Unknown target: {}", target); process::exit(1); }
}
println!("Generated {} DFA monitor app in {}/", target, output_dir);
println!(" Channels: {}", channels.len());
println!(" Window: {}", window);
println!(" Threshold: {:.3}", threshold);
if let Some(ref path) = db_path {
println!(" Source: {} (ogma-compatible db.json)", path);
}
}
#[allow(dead_code)]
struct Channel {
name: String,
c_type: String,
topic: String,
field: String,
msg_type: String,
}
fn parse_db_json(path: &str) -> Vec<Channel> {
let content = fs::read_to_string(path).unwrap_or_else(|e| {
eprintln!("Error reading {}: {}", path, e);
process::exit(1);
});
let mut channels = Vec::new();
let mut in_inputs = false;
let mut cur_name = String::new();
let mut cur_type = String::new();
let mut cur_topic = String::new();
let mut cur_field = String::new();
for line in content.lines() {
let t = line.trim();
if t.contains("\"inputs\"") { in_inputs = true; }
if !in_inputs { continue; }
if t.contains("\"name\"") {
if let Some(v) = extract_json_string(t) { cur_name = v; }
}
if t.contains("\"type\"") && !t.contains("\"fromType\"") && !t.contains("\"toType\"") {
if let Some(v) = extract_json_string_after(t, "\"type\"") { cur_type = v; }
}
if t.contains("\"topic\"") {
if let Some(v) = extract_json_string(t) { cur_topic = v; }
}
if t.contains("\"field\"") && !t.contains("\"fromField\"") {
if let Some(v) = extract_json_string(t) { cur_field = v; }
}
if (t == "}" || t == "},") && !cur_name.is_empty() && !cur_topic.is_empty() {
let clean_topic = cur_topic.trim_end_matches("_MID").to_string();
channels.push(Channel {
name: cur_name.clone(),
c_type: if cur_type.is_empty() { "double".into() } else { cur_type.clone() },
topic: clean_topic.clone(),
field: if cur_field.is_empty() { "payload".into() } else { cur_field.clone() },
msg_type: format!("{}_msg_t", clean_topic.to_lowercase()),
});
cur_name.clear();
cur_type.clear();
cur_topic.clear();
cur_field.clear();
}
if t.contains("\"topics\"") { in_inputs = false; }
}
channels
}
fn extract_json_string(line: &str) -> Option<String> {
let parts: Vec<&str> = line.split('"').collect();
if parts.len() >= 4 { Some(parts[3].to_string()) } else { None }
}
fn extract_json_string_after(line: &str, key: &str) -> Option<String> {
if let Some(pos) = line.find(key) {
let rest = &line[pos + key.len()..];
let parts: Vec<&str> = rest.split('"').collect();
if parts.len() >= 2 { return Some(parts[1].to_string()); }
}
None
}
fn generate_cfs_app_dir(out: &std::path::Path, channels: &[Channel], window: usize, threshold: f64) {
let src = out.join("fsw").join("src");
fs::create_dir_all(&src).unwrap();
fs::create_dir_all(out.join("fsw").join("platform_inc")).unwrap();
fs::create_dir_all(out.join("fsw").join("mission_inc")).unwrap();
fs::write(src.join("dfa_core.h"), generate_dfa_core_h()).unwrap();
fs::write(src.join("dfa_monitor_cfs.c"), generate_cfs_main(channels, window, threshold)).unwrap();
fs::write(src.join("dfa_monitor_cfs.h"), generate_cfs_header(channels)).unwrap();
fs::write(src.join("dfa_monitor_cfs_events.h"), DFA_EVENTS_H).unwrap();
fs::write(out.join("fsw").join("platform_inc").join("dfa_monitor_cfs_msgids.h"), generate_cfs_msgids(channels)).unwrap();
fs::write(out.join("fsw").join("mission_inc").join("dfa_monitor_cfs_perfids.h"), DFA_PERFIDS_H).unwrap();
fs::write(out.join("CMakeLists.txt"), CFS_CMAKE).unwrap();
}
fn generate_fprime_dir(out: &std::path::Path, channels: &[Channel], window: usize, threshold: f64) {
fs::create_dir_all(out).unwrap();
fs::write(out.join("dfa_core.h"), generate_dfa_core_h()).unwrap();
fs::write(out.join("DfaMonitor.fpp"), generate_fprime_fpp(channels)).unwrap();
fs::write(out.join("DfaMonitor.cpp"), generate_fprime_cpp(channels, window, threshold)).unwrap();
fs::write(out.join("DfaMonitor.hpp"), generate_fprime_hpp(channels, window)).unwrap();
fs::write(out.join("CMakeLists.txt"), FPRIME_CMAKE).unwrap();
}
fn generate_dfa_core_h() -> String {
let mut s = String::with_capacity(2048);
s.push_str("/* dfa_core.h -- DFA scaling analysis (compute only)\n");
s.push_str(" * Generated by struktura. Zero dependencies beyond <math.h>.\n");
s.push_str(" * Peng et al., Physical Review E 49(2), 1994.\n");
s.push_str(" * https://github.com/koscak-labs/struktura\n */\n\n");
s.push_str("#ifndef DFA_CORE_H\n#define DFA_CORE_H\n\n");
s.push_str("#include <math.h>\n\n");
s.push_str("#define DFA_NUM_BOXES 6\n");
s.push_str("static const int DFA_BOXES[DFA_NUM_BOXES] = {16, 24, 36, 54, 81, 121};\n\n");
s.push_str("typedef struct { double alpha; double r_squared; } dfa_result_t;\n\n");
s.push_str("static dfa_result_t dfa_compute(const double *values, int n) {\n");
s.push_str(" dfa_result_t result = {0.5, 0.0};\n");
s.push_str(" if (n < 64) return result;\n");
s.push_str(" double mean = 0.0;\n int i, seg, b;\n");
s.push_str(" for (i = 0; i < n; i++) mean += values[i];\n");
s.push_str(" mean /= (double)n;\n");
s.push_str(" double y[512];\n double cum = 0.0;\n");
s.push_str(" int nn = n > 512 ? 512 : n;\n");
s.push_str(" for (i = 0; i < nn; i++) { cum += values[i] - mean; y[i] = cum; }\n");
s.push_str(" double log_s[DFA_NUM_BOXES], log_f[DFA_NUM_BOXES];\n int pts = 0;\n");
s.push_str(" for (b = 0; b < DFA_NUM_BOXES && DFA_BOXES[b] <= nn / 4; b++) {\n");
s.push_str(" int s = DFA_BOXES[b], num_segs = nn / s;\n");
s.push_str(" if (num_segs == 0) continue;\n double f2_sum = 0.0;\n");
s.push_str(" for (seg = 0; seg < num_segs; seg++) {\n");
s.push_str(" int start = seg * s;\n");
s.push_str(" double sx=0,sy=0,sxy=0,sx2=0;\n");
s.push_str(" for (i = 0; i < s; i++) {\n");
s.push_str(" double xi = (double)i;\n");
s.push_str(" sx+=xi; sy+=y[start+i]; sxy+=xi*y[start+i]; sx2+=xi*xi;\n");
s.push_str(" }\n");
s.push_str(" double k=(double)s, det=k*sx2-sx*sx;\n");
s.push_str(" if (fabs(det)<1e-15) continue;\n");
s.push_str(" double a0=(sx2*sy-sx*sxy)/det, a1=(k*sxy-sx*sy)/det;\n");
s.push_str(" double resid=0.0;\n");
s.push_str(" for (i=0;i<s;i++){double d=y[start+i]-(a0+a1*(double)i);resid+=d*d;}\n");
s.push_str(" f2_sum += resid / k;\n }\n");
s.push_str(" double f = sqrt(f2_sum / (double)num_segs);\n");
s.push_str(" if (f > 0.0) { log_s[pts]=log((double)s); log_f[pts]=log(f); pts++; }\n");
s.push_str(" }\n if (pts < 3) return result;\n");
s.push_str(" double k=(double)pts, sx=0,sy=0,sxy=0,sx2=0;\n");
s.push_str(" for (i=0;i<pts;i++){sx+=log_s[i];sy+=log_f[i];sxy+=log_s[i]*log_f[i];sx2+=log_s[i]*log_s[i];}\n");
s.push_str(" double slope=(k*sxy-sx*sy)/(k*sx2-sx*sx);\n");
s.push_str(" double ic=(sy-slope*sx)/k, ym=sy/k, sst=0,ssr=0;\n");
s.push_str(" for (i=0;i<pts;i++){sst+=(log_f[i]-ym)*(log_f[i]-ym);ssr+=(log_f[i]-slope*log_s[i]-ic)*(log_f[i]-slope*log_s[i]-ic);}\n");
s.push_str(" result.alpha=slope; result.r_squared=1.0-ssr/(sst>1e-15?sst:1e-15);\n");
s.push_str(" return result;\n}\n\n#endif /* DFA_CORE_H */\n");
s
}
fn generate_cfs_main(channels: &[Channel], window: usize, threshold: f64) -> String {
let mut s = String::with_capacity(4096);
s.push_str("/* dfa_monitor_cfs.c -- Generated by struktura generate --cfs\n");
s.push_str(" * DFA structural health monitor for NASA cFS.\n");
s.push_str(" * https://github.com/koscak-labs/struktura\n */\n\n");
s.push_str("#include \"dfa_monitor_cfs.h\"\n");
s.push_str("#include \"dfa_monitor_cfs_events.h\"\n");
s.push_str("#include \"dfa_monitor_cfs_msgids.h\"\n");
s.push_str("#include \"dfa_core.h\"\n\n");
s.push_str(&format!("#define DFA_WINDOW_SIZE {}\n", window));
s.push_str(&format!("#define DFA_THRESHOLD {:.4}\n", threshold));
s.push_str("#define DFA_LEARN_WINDOWS 10\n");
s.push_str("#define DFA_R2_MIN 0.7\n\n");
s.push_str("typedef struct {\n");
s.push_str(&format!(" double buffer[{}];\n", window));
s.push_str(" uint32 pos;\n uint32 filled;\n");
s.push_str(" double baseline_alpha;\n uint8 baseline_set;\n uint32 window_count;\n");
s.push_str("} dfa_channel_t;\n\n");
for ch in channels {
s.push_str(&format!("static dfa_channel_t dfa_ch_{};\n", ch.name));
}
s.push_str("\nstatic CFE_SB_PipeId_t DFA_MONITOR_Pipe;\n");
s.push_str("static CFE_SB_Buffer_t *SBBufPtr;\n\n");
s.push_str("void DFA_MONITOR_AppMain(void) {\n");
s.push_str(" CFE_Status_t status;\n uint32 RunStatus = CFE_ES_RunStatus_APP_RUN;\n");
s.push_str(" status = DFA_MONITOR_Init();\n");
s.push_str(" if (status != CFE_SUCCESS) RunStatus = CFE_ES_RunStatus_APP_ERROR;\n");
s.push_str(" while (CFE_ES_RunLoop(&RunStatus) == true) {\n");
s.push_str(" status = CFE_SB_ReceiveBuffer(&SBBufPtr, DFA_MONITOR_Pipe, 500);\n");
s.push_str(" if (status == CFE_SUCCESS) DFA_MONITOR_ProcessPkt();\n");
s.push_str(" }\n CFE_ES_ExitApp(RunStatus);\n}\n\n");
s.push_str("CFE_Status_t DFA_MONITOR_Init(void) {\n");
s.push_str(" CFE_Status_t status;\n");
s.push_str(" CFE_EVS_Register(NULL, 0, CFE_EVS_EventFilter_BINARY);\n");
s.push_str(" status = CFE_SB_CreatePipe(&DFA_MONITOR_Pipe, 32, \"DFA_MON_PIPE\");\n");
s.push_str(" if (status != CFE_SUCCESS) return status;\n\n");
for ch in channels {
s.push_str(&format!(" CFE_SB_Subscribe(CFE_SB_ValueToMsgId({}_MID), DFA_MONITOR_Pipe);\n", ch.topic.to_uppercase()));
s.push_str(&format!(" memset(&dfa_ch_{}, 0, sizeof(dfa_ch_{}));\n", ch.name, ch.name));
}
s.push_str("\n CFE_EVS_SendEvent(DFA_MON_INIT_EID, CFE_EVS_EventType_INFORMATION,\n");
s.push_str(&format!(" \"DFA Monitor: {} channels, window={}, threshold={:.3}\");\n", channels.len(), window, threshold));
s.push_str(" return CFE_SUCCESS;\n}\n\n");
s.push_str("static void dfa_push(dfa_channel_t *ch, double value, const char *name) {\n");
s.push_str(&format!(" ch->buffer[ch->pos] = value;\n ch->pos = (ch->pos + 1) % {};\n", window));
s.push_str(" if (ch->pos == 0) ch->filled = 1;\n if (!ch->filled) return;\n");
s.push_str(" ch->window_count++;\n");
s.push_str(&format!(" dfa_result_t r = dfa_compute(ch->buffer, {});\n", window));
s.push_str(" if (!ch->baseline_set && ch->window_count >= DFA_LEARN_WINDOWS && r.r_squared > DFA_R2_MIN) {\n");
s.push_str(" ch->baseline_alpha = r.alpha;\n ch->baseline_set = 1;\n");
s.push_str(" CFE_EVS_SendEvent(DFA_MON_BASELINE_EID, CFE_EVS_EventType_INFORMATION,\n");
s.push_str(" \"DFA baseline %s: alpha=%.3f R2=%.4f\", name, r.alpha, r.r_squared);\n");
s.push_str(" return;\n }\n if (!ch->baseline_set) return;\n");
s.push_str(" if (r.r_squared < DFA_R2_MIN) return;\n");
s.push_str(" double shift = r.alpha - ch->baseline_alpha;\n");
s.push_str(" if (shift < 0) shift = -shift;\n");
s.push_str(" if (shift >= DFA_THRESHOLD) {\n");
s.push_str(" CFE_EVS_SendEvent(DFA_MON_SHIFT_EID, CFE_EVS_EventType_ERROR,\n");
s.push_str(" \"DFA SHIFT %s: alpha=%.3f baseline=%.3f delta=%.3f\",\n");
s.push_str(" name, r.alpha, ch->baseline_alpha, r.alpha - ch->baseline_alpha);\n");
s.push_str(" }\n}\n\n");
s.push_str("void DFA_MONITOR_ProcessPkt(void) {\n");
s.push_str(" CFE_SB_MsgId_t MsgId = CFE_SB_INVALID_MSG_ID;\n");
s.push_str(" CFE_MSG_GetMsgId(&SBBufPtr->Msg, &MsgId);\n\n");
for (i, ch) in channels.iter().enumerate() {
let kw = if i == 0 { "if" } else { "else if" };
s.push_str(&format!(" {} (CFE_SB_MsgId_Equal(MsgId, CFE_SB_ValueToMsgId({}_MID))) {{\n", kw, ch.topic.to_uppercase()));
s.push_str(&format!(" {} *msg = ({}*)&SBBufPtr->Msg;\n", ch.msg_type, ch.msg_type));
s.push_str(&format!(" dfa_push(&dfa_ch_{}, (double)msg->{}, \"{}\");\n", ch.name, ch.field, ch.name));
s.push_str(" }\n");
}
s.push_str("}\n");
s
}
fn generate_cfs_header(channels: &[Channel]) -> String {
let mut s = String::new();
s.push_str("#ifndef DFA_MONITOR_CFS_H\n#define DFA_MONITOR_CFS_H\n\n");
s.push_str("#include \"cfe.h\"\n#include <string.h>\n#include <math.h>\n\n");
s.push_str("void DFA_MONITOR_AppMain(void);\n");
s.push_str("CFE_Status_t DFA_MONITOR_Init(void);\n");
s.push_str("void DFA_MONITOR_ProcessPkt(void);\n\n");
let _ = channels;
s.push_str("#endif\n");
s
}
fn generate_cfs_msgids(channels: &[Channel]) -> String {
let mut s = String::new();
s.push_str("#ifndef DFA_MONITOR_CFS_MSGIDS_H\n#define DFA_MONITOR_CFS_MSGIDS_H\n\n");
for (i, ch) in channels.iter().enumerate() {
s.push_str(&format!("#define {}_MID 0x{:04X}\n", ch.topic.to_uppercase(), 0x1900 + i));
}
s.push_str("\n#endif\n");
s
}
fn generate_fprime_fpp(channels: &[Channel]) -> String {
let mut s = String::with_capacity(2048);
s.push_str("module Svc {\n");
s.push_str(" @ DFA structural health monitor\n");
s.push_str(" @ Generated by: struktura generate --fprime\n");
s.push_str(" @ https://github.com/koscak-labs/struktura\n");
s.push_str(" passive component DfaMonitor {\n\n");
s.push_str(" sync input port schedIn: Svc.Sched\n\n");
for ch in channels {
s.push_str(&format!(" guarded input port {}In: Fw.Tlm\n", ch.name));
}
s.push_str("\n event StructuralShift(\n");
s.push_str(" channelName: string size 32\n");
s.push_str(" baseline_alpha: F64\n current_alpha: F64\n delta: F64\n");
s.push_str(" ) severity warning high\n\n");
s.push_str(" event BaselineEstablished(\n");
s.push_str(" channelName: string size 32\n alpha: F64\n r_squared: F64\n");
s.push_str(" ) severity activity high\n\n");
s.push_str(" telemetry DfaAlpha: F64\n telemetry DfaR2: F64\n\n");
s.push_str(" time get port timeCaller\n event port logOut\n telemetry port tlmOut\n");
s.push_str(" }\n}\n");
s
}
fn generate_fprime_cpp(channels: &[Channel], window: usize, threshold: f64) -> String {
let mut s = String::with_capacity(2048);
s.push_str("// DfaMonitor.cpp -- Generated by struktura generate --fprime\n");
s.push_str("// https://github.com/koscak-labs/struktura\n\n");
s.push_str("#include \"DfaMonitor.hpp\"\n#include \"dfa_core.h\"\n\n");
s.push_str("namespace Svc {\n\n");
s.push_str("DfaMonitor::DfaMonitor(const char* name) : DfaMonitorComponentBase(name) {\n");
for ch in channels {
s.push_str(&format!(" memset(&m_ch_{}, 0, sizeof(m_ch_{}));\n", ch.name, ch.name));
}
s.push_str("}\n\n");
for ch in channels {
s.push_str(&format!("void DfaMonitor::{}In_handler(NATIVE_INT_TYPE portNum, FwTlmBuffer& val) {{\n", ch.name));
s.push_str(" F64 v; val.deserialize(v);\n");
s.push_str(&format!(" pushSample(m_ch_{}, v, \"{}\");\n", ch.name, ch.name));
s.push_str("}\n\n");
}
s.push_str("void DfaMonitor::pushSample(DfaChannel& ch, F64 value, const char* name) {\n");
s.push_str(&format!(" ch.buffer[ch.pos] = value;\n ch.pos = (ch.pos + 1) % {};\n", window));
s.push_str(" if (ch.pos == 0) ch.filled = true;\n if (!ch.filled) return;\n");
s.push_str(" ch.windowCount++;\n");
s.push_str(&format!(" dfa_result_t r = dfa_compute(ch.buffer, {});\n", window));
s.push_str(" if (!ch.baselineSet && ch.windowCount >= 10 && r.r_squared > 0.7) {\n");
s.push_str(" ch.baselineAlpha = r.alpha; ch.baselineSet = true;\n");
s.push_str(" this->log_ACTIVITY_HI_BaselineEstablished(name, r.alpha, r.r_squared);\n");
s.push_str(" return;\n }\n if (!ch.baselineSet || r.r_squared < 0.7) return;\n");
s.push_str(" F64 shift = (r.alpha > ch.baselineAlpha) ? r.alpha - ch.baselineAlpha : ch.baselineAlpha - r.alpha;\n");
s.push_str(&format!(" if (shift >= {:.4}) {{\n", threshold));
s.push_str(" this->log_WARNING_HI_StructuralShift(name, ch.baselineAlpha, r.alpha, r.alpha - ch.baselineAlpha);\n");
s.push_str(" }\n this->tlmWrite_DfaAlpha(r.alpha);\n this->tlmWrite_DfaR2(r.r_squared);\n");
s.push_str("}\n\n} // namespace Svc\n");
s
}
fn generate_fprime_hpp(channels: &[Channel], window: usize) -> String {
let mut s = String::with_capacity(1024);
s.push_str("#ifndef DFA_MONITOR_HPP\n#define DFA_MONITOR_HPP\n\n");
s.push_str("#include \"DfaMonitorComponentAc.hpp\"\n\n");
s.push_str("namespace Svc {\n\n");
s.push_str("struct DfaChannel {\n");
s.push_str(&format!(" double buffer[{}];\n", window));
s.push_str(" U32 pos; bool filled; double baselineAlpha;\n");
s.push_str(" bool baselineSet; U32 windowCount;\n};\n\n");
s.push_str("class DfaMonitor : public DfaMonitorComponentBase {\n");
s.push_str(" public:\n DfaMonitor(const char* name);\n");
s.push_str(" private:\n");
for ch in channels {
s.push_str(&format!(" void {}In_handler(NATIVE_INT_TYPE portNum, FwTlmBuffer& val);\n", ch.name));
}
s.push_str(" void pushSample(DfaChannel& ch, F64 value, const char* name);\n");
for ch in channels {
s.push_str(&format!(" DfaChannel m_ch_{};\n", ch.name));
}
s.push_str("};\n\n} // namespace Svc\n\n#endif\n");
s
}
const DFA_EVENTS_H: &str = "\
#ifndef DFA_MONITOR_CFS_EVENTS_H
#define DFA_MONITOR_CFS_EVENTS_H
#define DFA_MON_INIT_EID 1
#define DFA_MON_BASELINE_EID 2
#define DFA_MON_SHIFT_EID 3
#endif
";
const DFA_PERFIDS_H: &str = "\
#ifndef DFA_MONITOR_CFS_PERFIDS_H
#define DFA_MONITOR_CFS_PERFIDS_H
#define DFA_MON_PERF_ID 91
#endif
";
const CFS_CMAKE: &str = "\
cmake_minimum_required(VERSION 2.6.4)
project(CFE_DFA_MONITOR_APP C)
include_directories(fsw/mission_inc)
include_directories(fsw/platform_inc)
aux_source_directory(fsw/src APP_SRC_FILES)
add_cfe_app(dfa_monitor_cfs ${APP_SRC_FILES})
";
const FPRIME_CMAKE: &str = "\
set(SOURCE_FILES DfaMonitor.cpp)
set(MOD_DEPS Fw/Tlm Svc/Sched)
register_fprime_module()
";
fn generate_ros_dir(out: &std::path::Path, channels: &[Channel], window: usize, threshold: f64) {
let src = out.join("src");
fs::create_dir_all(&src).unwrap();
fs::write(src.join("dfa_core.h"), generate_dfa_core_h()).unwrap();
fs::write(src.join("dfa_monitor_node.cpp"), generate_ros_monitor(channels, window, threshold)).unwrap();
fs::write(out.join("CMakeLists.txt"), generate_ros_cmake()).unwrap();
fs::write(out.join("package.xml"), generate_ros_package()).unwrap();
}
fn generate_ros_monitor(channels: &[Channel], window: usize, threshold: f64) -> String {
let mut s = String::with_capacity(4096);
s.push_str("/* dfa_monitor_node.cpp -- Generated by struktura generate --ros\n");
s.push_str(" * DFA structural health monitor for ROS 2.\n");
s.push_str(" * https://github.com/koscak-labs/struktura\n */\n\n");
s.push_str("#include <functional>\n#include <memory>\n#include <cmath>\n#include <cstring>\n\n");
s.push_str("#include \"rclcpp/rclcpp.hpp\"\n");
s.push_str("#include \"std_msgs/msg/float64.hpp\"\n");
s.push_str("#include \"std_msgs/msg/string.hpp\"\n\n");
s.push_str("#include \"dfa_core.h\"\n\n");
s.push_str(&format!("#define DFA_WINDOW_SIZE {}\n", window));
s.push_str(&format!("#define DFA_THRESHOLD {:.4}\n", threshold));
s.push_str("#define DFA_LEARN_WINDOWS 10\n#define DFA_R2_MIN 0.7\n\n");
s.push_str("using std::placeholders::_1;\n\n");
s.push_str("struct DfaChannel {\n");
s.push_str(&format!(" double buffer[{}];\n", window));
s.push_str(" uint32_t pos = 0;\n bool filled = false;\n");
s.push_str(" double baseline_alpha = 0.0;\n bool baseline_set = false;\n");
s.push_str(" uint32_t window_count = 0;\n};\n\n");
s.push_str("class DfaMonitorNode : public rclcpp::Node {\n");
s.push_str("public:\n");
s.push_str(" DfaMonitorNode() : Node(\"dfa_monitor\") {\n");
for ch in channels {
s.push_str(&format!(" {}_sub_ = this->create_subscription<std_msgs::msg::Float64>(\n", ch.name));
s.push_str(&format!(" \"/{}\", 10,\n", ch.name));
s.push_str(&format!(" std::bind(&DfaMonitorNode::{}_callback, this, _1));\n", ch.name));
s.push_str(&format!(" std::memset(&ch_{}, 0, sizeof(ch_{}));\n\n", ch.name, ch.name));
}
s.push_str(" shift_pub_ = this->create_publisher<std_msgs::msg::String>(\"dfa/shift\", 10);\n");
s.push_str(&format!(" RCLCPP_INFO(this->get_logger(), \"DFA Monitor: {} channels, window={}, threshold={:.3}\");\n", channels.len(), window, threshold));
s.push_str(" }\n\nprivate:\n");
s.push_str(" void push_sample(DfaChannel &ch, double value, const char *name) {\n");
s.push_str(&format!(" ch.buffer[ch.pos] = value;\n ch.pos = (ch.pos + 1) % {};\n", window));
s.push_str(" if (ch.pos == 0) ch.filled = true;\n if (!ch.filled) return;\n");
s.push_str(" ch.window_count++;\n");
s.push_str(&format!(" dfa_result_t r = dfa_compute(ch.buffer, {});\n", window));
s.push_str(" if (!ch.baseline_set && ch.window_count >= DFA_LEARN_WINDOWS && r.r_squared > DFA_R2_MIN) {\n");
s.push_str(" ch.baseline_alpha = r.alpha; ch.baseline_set = true;\n");
s.push_str(" RCLCPP_INFO(this->get_logger(), \"DFA baseline %s: alpha=%.3f R2=%.4f\", name, r.alpha, r.r_squared);\n");
s.push_str(" return;\n }\n if (!ch.baseline_set || r.r_squared < DFA_R2_MIN) return;\n");
s.push_str(" double shift = std::abs(r.alpha - ch.baseline_alpha);\n");
s.push_str(" if (shift >= DFA_THRESHOLD) {\n");
s.push_str(" auto msg = std_msgs::msg::String();\n");
s.push_str(" char buf[128];\n");
s.push_str(" std::snprintf(buf, sizeof(buf), \"SHIFT %s: alpha=%.3f baseline=%.3f delta=%.3f\",\n");
s.push_str(" name, r.alpha, ch.baseline_alpha, r.alpha - ch.baseline_alpha);\n");
s.push_str(" msg.data = buf;\n shift_pub_->publish(msg);\n");
s.push_str(" RCLCPP_WARN(this->get_logger(), \"%s\", buf);\n");
s.push_str(" }\n }\n\n");
for ch in channels {
s.push_str(&format!(" void {}_callback(const std_msgs::msg::Float64::SharedPtr msg) {{\n", ch.name));
s.push_str(&format!(" push_sample(ch_{}, msg->data, \"{}\");\n", ch.name, ch.name));
s.push_str(" }\n\n");
}
for ch in channels {
s.push_str(&format!(" rclcpp::Subscription<std_msgs::msg::Float64>::SharedPtr {}_sub_;\n", ch.name));
s.push_str(&format!(" DfaChannel ch_{};\n", ch.name));
}
s.push_str(" rclcpp::Publisher<std_msgs::msg::String>::SharedPtr shift_pub_;\n");
s.push_str("};\n\n");
s.push_str("int main(int argc, char *argv[]) {\n");
s.push_str(" rclcpp::init(argc, argv);\n");
s.push_str(" rclcpp::spin(std::make_shared<DfaMonitorNode>());\n");
s.push_str(" rclcpp::shutdown();\n return 0;\n}\n");
s
}
fn generate_ros_cmake() -> String {
let mut s = String::new();
s.push_str("cmake_minimum_required(VERSION 3.8)\nproject(dfa_monitor)\n\n");
s.push_str("find_package(ament_cmake REQUIRED)\nfind_package(rclcpp REQUIRED)\n");
s.push_str("find_package(std_msgs REQUIRED)\n\n");
s.push_str("add_executable(dfa_monitor_node src/dfa_monitor_node.cpp)\n");
s.push_str("ament_target_dependencies(dfa_monitor_node rclcpp std_msgs)\n\n");
s.push_str("install(TARGETS dfa_monitor_node DESTINATION lib/${PROJECT_NAME})\n");
s.push_str("ament_package()\n");
s
}
fn cmd_benchmark_faults() {
use struktura::dfa;
use struktura::space::synth_structural_fault;
let n = 4096;
let n_seeds = 20u64;
fn lcg_next(state: &mut u64) -> f64 {
*state = state.wrapping_mul(6364136223846793005).wrapping_add(1442695040888963407);
(*state >> 33) as f64 / (1u64 << 31) as f64 - 0.5
}
let fault_at = n / 2;
type Injector = fn(&[f64], u64, usize) -> Vec<f64>;
let inject_structural: Injector = |_clean, seed, _fa| synth_structural_fault(4096, seed, 0.5);
let inject_spike: Injector = |clean, seed, fa| {
let mut v = clean.to_vec();
let mut s = seed + 10;
for (i, val) in v.iter_mut().enumerate().skip(fa) {
if i % 50 == 0 { *val += lcg_next(&mut s) * 0.5; }
}
v
};
let inject_stuck: Injector = |clean, _seed, fa| {
let mut v = clean.to_vec();
let stuck_val = v[fa];
let end = (fa + 200).min(v.len());
for val in v[fa..end].iter_mut() { *val = stuck_val; }
v
};
let inject_drift: Injector = |clean, _seed, fa| {
let mut v = clean.to_vec();
for (i, val) in v.iter_mut().enumerate().skip(fa) { *val += (i - fa) as f64 * 0.0001; }
v
};
let inject_regime: Injector = |clean, _seed, fa| {
let mut v = clean.to_vec();
for val in v.iter_mut().skip(fa) { *val += 0.15; }
v
};
let inject_packet_loss: Injector = |clean, seed, fa| {
let mut v = clean.to_vec();
let mut s = seed + 20;
for val in v.iter_mut().skip(fa) {
if lcg_next(&mut s).abs() < 0.1 { *val = 0.0; }
}
v
};
let fault_types: Vec<(&str, &str, Injector)> = vec![
("structural", "correlation", inject_structural),
("spike", "value", inject_spike),
("stuck", "value", inject_stuck),
("drift", "value", inject_drift),
("regime_shift", "value", inject_regime),
("packet_loss", "value", inject_packet_loss),
];
println!();
println!(" \x1b[1mSTRUKTURA FAULT TYPE BENCHMARK\x1b[0m");
println!(" Statistical detection test: {} seeds per fault type", n_seeds);
println!(" ================================================================");
println!();
let mut null_shifts: Vec<f64> = Vec::new();
for seed in 1..=n_seeds {
let clean = synth_structural_fault(n, seed * 7919, 1.1);
let a1 = dfa(&clean[..fault_at]).alpha;
let a2 = dfa(&clean[fault_at..]).alpha;
null_shifts.push((a2 - a1).abs());
}
let mut sorted_null = null_shifts.clone();
sorted_null.sort_by(|a, b| a.partial_cmp(b).unwrap());
let p95_idx = ((sorted_null.len() as f64 * 0.95) as usize).min(sorted_null.len() - 1);
let p95 = sorted_null[p95_idx];
let null_mean = null_shifts.iter().sum::<f64>() / null_shifts.len() as f64;
println!(" Null distribution (clean signal, {} seeds):", n_seeds);
println!(" mean |Δα| = {:.4}, 95th percentile = {:.4}", null_mean, p95);
println!(" Detection = fault shift > {:.4} (95th pct of clean shifts)", p95);
println!();
println!(" | Fault Type | Class | Mean |Δα| | vs Null | Detect Rate |");
println!(" |----------------|-------------|-----------|---------|-------------|");
for (name, class, inject) in &fault_types {
let mut shifts = Vec::new();
let mut detections = 0usize;
for seed in 1..=n_seeds {
let s = seed * 7919;
let clean = synth_structural_fault(n, s, 1.1);
let baseline_alpha = dfa(&clean[..fault_at]).alpha;
let faulted = inject(&clean, s, fault_at);
let fault_alpha = dfa(&faulted[fault_at..]).alpha;
let shift = (fault_alpha - baseline_alpha).abs();
shifts.push(shift);
if shift > p95 { detections += 1; }
}
let mean_shift = shifts.iter().sum::<f64>() / shifts.len() as f64;
let ratio = if null_mean > 0.0001 { mean_shift / null_mean } else { 0.0 };
let rate = detections as f64 / n_seeds as f64;
let rate_color = if rate >= 0.8 { "\x1b[32m" } else if rate >= 0.4 { "\x1b[33m" } else { "\x1b[31m" };
println!(
" | {:<14} | {:<11} | {:.4} | {:.1}x | {}{:>3.0}%\x1b[0m |",
name, class, mean_shift, ratio, rate_color, rate * 100.0
);
}
println!();
println!(" Detect Rate = fraction of seeds where fault shift exceeds the");
println!(" 95th percentile of clean-signal shifts (false positive rate ≤ 5%).");
println!();
println!(" DFA measures STRUCTURE (long-range correlations), not VALUES.");
println!(" Residual-based detectors (xLSTM, ARIMA) measure VALUE deviations.");
println!(" Together they cover the full fault taxonomy, neither alone does.");
println!();
println!(" Algorithm: Detrended Fluctuation Analysis (Peng 1994)");
println!(" Signal: {} samples/seed, fault injected at sample {}", n, fault_at);
println!(" https://crates.io/crates/struktura");
println!();
}
fn cmd_benchmark_telemetry(args: &[String]) {
use struktura::telemetry_bench::{run_benchmark, CHANNEL_NAMES};
let mut length = 700usize;
let mut n_seeds = 200u64;
let mut n_null = 200u64;
let mut i = 2;
while i < args.len() {
match args[i].as_str() {
"--length" if i + 1 < args.len() => { length = args[i + 1].parse().unwrap_or(700); i += 2; }
"--seeds" if i + 1 < args.len() => { n_seeds = args[i + 1].parse().unwrap_or(200); i += 2; }
"--null-seeds" if i + 1 < args.len() => { n_null = args[i + 1].parse().unwrap_or(200); i += 2; }
_ => { i += 1; }
}
}
println!();
println!(" \x1b[1mSTRUKTURA COUPLED-TELEMETRY BENCHMARK\x1b[0m");
println!(" 6-channel spacecraft sim (power/thermal/wheel/pointing/payload)");
println!(" Fault window: 58%..70% of sequence | {} samples", length);
println!(" {} null seeds (threshold calibration) + {} eval seeds (disjoint)", n_null, n_seeds);
println!(" ================================================================");
println!();
let report = run_benchmark(length, n_seeds, n_null);
println!(" Per-channel Bonferroni threshold |Δα| (99.17th pct of null, family FPR <= 5%):");
for (ch, p) in report.thresholds.iter().enumerate() {
println!(" {:<16} {:.4}", CHANNEL_NAMES[ch], p);
}
println!();
println!(" \x1b[1mEmpirical family-wise FPR on {} held-out clean pairs: {:.1}%\x1b[0m",
n_seeds, report.empirical_fpr * 100.0);
println!();
println!(" | Fault Type | Detect Rate | Mean Max |Δα| | Best Channel |");
println!(" |--------------------|-------------|---------------|------------------|");
for r in &report.results {
let rate_color = if r.detect_rate >= 0.8 { "\x1b[32m" }
else if r.detect_rate >= 0.4 { "\x1b[33m" }
else { "\x1b[31m" };
println!(
" | {:<18} | {}{:>4.0}%\x1b[0m | {:.4} | {:<16} |",
r.fault, rate_color, r.detect_rate * 100.0, r.mean_max_shift,
CHANNEL_NAMES[r.best_channel]
);
}
println!();
println!(" Detection = any channel's |Δα| (calibration vs test) exceeds that");
println!(" channel's Bonferroni-corrected null threshold. A detect rate is only");
println!(" meaningful relative to the measured FPR line above.");
println!(" correlation_change = structural fault the additive taxonomy misses.");
println!();
use struktura::telemetry_bench::correlation_resolution_curve;
let fracs = [0.12, 0.20, 0.30, 0.42];
let curve = correlation_resolution_curve(length, n_seeds, n_null, &fracs);
println!(" \x1b[1mSTRUCTURAL FAULT RESOLUTION CURVE\x1b[0m (correlation_change)");
println!(" | Fault duration | Samples | Detect Rate |");
println!(" |----------------|---------|-------------|");
for (frac, rate) in &curve {
let samples = (length as f64 * frac) as usize;
let color = if *rate >= 0.8 { "\x1b[32m" } else if *rate >= 0.4 { "\x1b[33m" } else { "\x1b[31m" };
println!(" | {:>4.0}% of seq | {:<7} | {}{:>4.0}%\x1b[0m |",
frac * 100.0, samples, color, rate * 100.0);
}
println!();
use struktura::telemetry_bench::timestep_f1_benchmark;
let f1_seeds = n_seeds.min(30);
let f1 = timestep_f1_benchmark(length, f1_seeds, 96, 2);
println!(" \x1b[1mPER-TIMESTEP F1\x1b[0m (trailing DFA window 96, step 2, {} seeds,", f1_seeds);
println!(" threshold = 99th pct of calibration scores, residual-detector protocol)");
println!(" | Fault Type | F1 (step) | False Alarm | Event Detect | Latency |");
println!(" |--------------------|-----------|-------------|--------------|---------|");
for r in &f1 {
let color = if r.f1 >= 0.5 { "\x1b[32m" } else if r.f1 >= 0.2 { "\x1b[33m" } else { "\x1b[31m" };
let ev_color = if r.event_detect_rate >= 0.8 { "\x1b[32m" }
else if r.event_detect_rate >= 0.4 { "\x1b[33m" } else { "\x1b[31m" };
println!(
" | {:<18} | {}{:.3}\x1b[0m | {:.4} | {}{:>4.0}%\x1b[0m | {:>4.0} |",
r.fault, color, r.f1, r.false_alarm_rate,
ev_color, r.event_detect_rate * 100.0, r.mean_latency
);
}
println!();
println!(" F1 (step) = per-timestep, punishes latency. Event Detect = flag lands");
println!(" within [fault_start, fault_end + window] (NAB-style event scoring).");
println!(" Latency = mean samples from fault start to first flag.");
println!();
use struktura::telemetry_bench::hybrid_benchmark;
let (hybrid, clean_far) = hybrid_benchmark(length, f1_seeds, 96, 2);
println!(" \x1b[1mHYBRID MONITOR\x1b[0m (AR1-residual + repeated-value + DFA + CUSUM, OR-fused,");
println!(" every threshold calibrated on the clean calibration sequence)");
println!(" Event-level false-alarm rate on {} clean sequences: {:.1}%", f1_seeds, clean_far * 100.0);
println!(" | Fault Type | Event Detect | Latency | First (res/rep/dfa/cusum) |");
println!(" |--------------------|--------------|---------|---------------------------|");
let mut covered = 0usize;
for r in &hybrid {
let color = if r.event_detect_rate >= 0.8 { "\x1b[32m" }
else if r.event_detect_rate >= 0.4 { "\x1b[33m" } else { "\x1b[31m" };
if r.event_detect_rate >= 0.8 { covered += 1; }
println!(
" | {:<18} | {}{:>4.0}%\x1b[0m | {:>4.0} | {}/{}/{}/{} |",
r.fault, color, r.event_detect_rate * 100.0, r.mean_latency,
r.first_detector.0, r.first_detector.1, r.first_detector.2, r.first_detector.3
);
}
println!();
println!(" \x1b[1mTaxonomy coverage: {}/7 fault types at >=80% event detection\x1b[0m", covered);
println!();
println!(" Algorithm: Detrended Fluctuation Analysis (Peng 1994)");
println!(" https://crates.io/crates/struktura");
println!();
}
fn cmd_monitor_perf() {
use struktura::monitor::HybridMonitor;
use struktura::telemetry_bench::synth_spacecraft;
use std::time::Instant;
println!();
println!(" \x1b[1mSTREAMING MONITOR PERFORMANCE\x1b[0m (flight-relevant timing)");
println!(" 6 channels, window 96, DFA stride 2, release build, this host");
println!();
let calib = synth_spacecraft(2048, 424242);
let t0 = Instant::now();
let mut mon = HybridMonitor::calibrate(&calib).expect("calibration");
let calib_time = t0.elapsed();
println!(" Calibration (2048 samples x 6 ch): {:?}", calib_time);
let n = 200_000usize;
let stream = synth_spacecraft(n, 434343);
let mut sample = [0.0f64; 6];
let mut worst_ns: u128 = 0;
let mut total_ns: u128 = 0;
let mut alarms = 0usize;
let mut leg_counts = [0usize; 5];
let mut lat_ns: Vec<u64> = Vec::with_capacity(n);
for t in 0..n {
for ch in 0..6 {
sample[ch] = stream[ch][t];
}
let s = Instant::now();
if let Some(leg) = mon.push(&sample) {
alarms += 1;
use struktura::monitor::Leg;
match leg {
Leg::Residual => leg_counts[0] += 1,
Leg::RepeatedValue => leg_counts[1] += 1,
Leg::Dfa => leg_counts[2] += 1,
Leg::LevelShift => leg_counts[3] += 1,
Leg::ResidualCusum => leg_counts[4] += 1,
Leg::Missingness => {}
Leg::Parity => {}
}
mon.reset();
}
let e = s.elapsed().as_nanos();
total_ns += e;
lat_ns.push(e as u64);
if e > worst_ns {
worst_ns = e;
}
}
lat_ns.sort_unstable();
let pct = |p: f64| lat_ns[((lat_ns.len() as f64 * p) as usize).min(lat_ns.len() - 1)];
let mean_ns = total_ns / n as u128;
println!(" Stream: {} samples x 6 channels", n);
println!(" Mean per-sample cost: {:>8} ns", mean_ns);
println!(" P50 / P99 / P99.9: {:>8} / {} / {} ns", pct(0.50), pct(0.99), pct(0.999));
println!(" Raw max (includes OS jitter): {:>6} ns", worst_ns);
println!(" Throughput: {:>8.1} Msamples/s", 1000.0 / mean_ns as f64);
println!(" Alarms on clean stream: {:>8} ({} samples ≈ {:.1}x calib length)",
alarms, n, n as f64 / 2048.0);
println!(" By leg (res/rep/dfa/level/cusum): {}/{}/{}/{}/{}",
leg_counts[0], leg_counts[1], leg_counts[2], leg_counts[3], leg_counts[4]);
println!();
println!(" Memory: fixed after calibration, per channel {} + {} f64 rings;",
struktura::monitor::WINDOW, struktura::monitor::ROLL);
println!(" no heap allocation in the push path.");
println!();
use struktura::monitor::classify_alarm;
use struktura::telemetry_bench::{inject_fault, inject_fault_validity, FAULT_TYPES};
let mut class_correct = [0usize; 7];
let mut class_total = [0usize; 7];
let mut class_confusion: Vec<(String, String, usize)> = Vec::new();
let seeds = 20u64;
let len = 2048usize;
let fstart = (len as f64 * 0.58) as usize;
let fstop = fstart + ((len as f64 * 0.12) as usize).max(8);
println!(" \x1b[1mFAULT DETECTION\x1b[0m (streaming monitor, {} seeds, {}-sample streams)", seeds, len);
println!(" | Fault Type | Detect | Mean Latency |");
println!(" |--------------------|--------|--------------|");
for fault in FAULT_TYPES.iter() {
let mut hits = 0usize;
let mut lat_sum = 0.0f64;
for seed in 1..=seeds {
let s = seed * 7919;
let cal = synth_spacecraft(len, s + 100);
let clean_test = synth_spacecraft(len, s + 200);
let faulted = inject_fault(&clean_test, fault, s);
let validity = inject_fault_validity(fault, len);
let mut m = match HybridMonitor::calibrate(&cal) { Some(m) => m, None => continue };
let mut smp = [0.0f64; 6];
let mut vld = [true; 6];
for t in 0..len {
for ch in 0..6 {
smp[ch] = faulted[ch][t];
vld[ch] = validity[ch][t];
}
if m.push_with_validity(&smp, &vld).is_some() {
if t >= fstart && t < fstop + 96 {
hits += 1;
lat_sum += (t - fstart) as f64;
if let Some(report) = m.last_alarm() {
let fi = FAULT_TYPES.iter().position(|f| f == fault).unwrap();
class_total[fi] += 1;
let report = if report.leg == struktura::monitor::Leg::Parity {
let mut d = HybridMonitor::calibrate(&cal).unwrap();
d.set_leg_enabled(struktura::monitor::Leg::Parity, false);
let mut smp2 = [0.0f64; 6];
let mut vld2 = [true; 6];
let mut rep2 = report;
for t2 in 0..len {
for ch in 0..6 {
smp2[ch] = faulted[ch][t2];
vld2[ch] = validity[ch][t2];
}
if d.push_with_validity(&smp2, &vld2).is_some() {
if let Some(r2) = d.last_alarm() {
rep2 = r2;
}
break;
}
}
rep2
} else {
report
};
let predicted = classify_alarm(&report);
let correct = predicted == *fault
|| (*fault == "mixed"
&& matches!(predicted, "packet_loss" | "spike" | "drift"));
if correct {
class_correct[fi] += 1;
} else if let Some(e) = class_confusion.iter_mut().find(
|(a, p, _)| a == fault && p == predicted,
) {
e.2 += 1;
} else {
class_confusion.push(
(fault.to_string(), predicted.to_string(), 1),
);
}
}
}
break;
}
}
}
let rate = hits as f64 / seeds as f64;
let color = if rate >= 0.8 { "\x1b[32m" } else if rate >= 0.4 { "\x1b[33m" } else { "\x1b[31m" };
println!(" | {:<18} | {}{:>4.0}%\x1b[0m | {:>6.0} |",
fault, color, rate * 100.0, if hits > 0 { lat_sum / hits as f64 } else { 0.0 });
}
println!();
println!(" \x1b[1mFAULT-CLASS IDENTIFICATION\x1b[0m (from alarm provenance, rule-based)");
println!(" | Fault Type | Correct ID |");
println!(" |--------------------|------------|");
for (fi, fault) in FAULT_TYPES.iter().enumerate() {
if class_total[fi] == 0 { continue; }
let acc = class_correct[fi] as f64 / class_total[fi] as f64;
let color = if acc >= 0.8 { "\x1b[32m" } else if acc >= 0.4 { "\x1b[33m" } else { "\x1b[31m" };
println!(" | {:<18} | {}{:>4.0}%\x1b[0m |", fault, color, acc * 100.0);
}
if !class_confusion.is_empty() {
println!(" Misclassifications:");
for (actual, predicted, count) in &class_confusion {
println!(" {} -> {} ({}x)", actual, predicted, count);
}
}
println!();
}
fn cmd_monitor_real() {
use struktura::monitor::HybridMonitor;
use std::path::Path;
println!();
println!(" \x1b[1mMONITOR ON REAL NASA DATA\x1b[0m");
println!(" ================================================================");
let ims_dir = Path::new("data/ims/extracted/2nd_test/2nd_test");
let cache = Path::new("data/ims_monitor_stream.csv");
let stream: Vec<[f64; 4]> = if cache.exists() {
std::fs::read_to_string(cache)
.unwrap_or_default()
.lines()
.filter_map(|l| {
let mut it = l.split(',');
let mut row = [0.0f64; 4];
for slot in row.iter_mut() {
*slot = it.next()?.trim().parse().ok()?;
}
Some(row)
})
.collect()
} else if ims_dir.exists() {
let mut names: Vec<_> = std::fs::read_dir(ims_dir)
.expect("read ims dir")
.filter_map(|e| e.ok().map(|e| e.path()))
.collect();
names.sort();
println!(" Computing per-recording RMS from {} raw recordings...", names.len());
let mut rows = Vec::with_capacity(names.len());
for p in &names {
let content = match std::fs::read_to_string(p) { Ok(c) => c, Err(_) => continue };
let mut sum_sq = [0.0f64; 4];
let mut n = 0usize;
for line in content.lines() {
let mut it = line.split_whitespace();
let mut vals = [0.0f64; 4];
let mut ok = true;
for slot in vals.iter_mut() {
match it.next().and_then(|s| s.parse::<f64>().ok()) {
Some(v) => *slot = v,
None => { ok = false; break; }
}
}
if ok {
for ch in 0..4 { sum_sq[ch] += vals[ch] * vals[ch]; }
n += 1;
}
}
if n > 0 {
let mut row = [0.0f64; 4];
for ch in 0..4 { row[ch] = (sum_sq[ch] / n as f64).sqrt(); }
rows.push(row);
}
}
let out: String = rows
.iter()
.map(|r| format!("{:.6},{:.6},{:.6},{:.6}\n", r[0], r[1], r[2], r[3]))
.collect();
let _ = std::fs::write(cache, out);
rows
} else {
Vec::new()
};
if stream.len() >= 400 {
let n = stream.len();
let calib_n = 300usize;
let calib: Vec<Vec<f64>> = (0..4)
.map(|ch| stream[..calib_n].iter().map(|r| r[ch]).collect())
.collect();
let calib_short: Vec<Vec<f64>> = (0..4)
.map(|ch| stream[..200].iter().map(|r| r[ch]).collect())
.collect();
if let Some(mut mon) = HybridMonitor::calibrate(&calib_short) {
let mut fa = None;
for (i, row) in stream[200..300].iter().enumerate() {
if let Some(leg) = mon.push(row) {
fa = Some((200 + i, leg));
break;
}
}
match fa {
Some((idx, leg)) => println!(
" \x1b[31mCONTROL FAILED: alarm on held-out healthy segment at {} via {:?}\x1b[0m",
idx, leg
),
None => println!(" Control PASS: no alarm on held-out healthy recordings 200..300"),
}
}
match HybridMonitor::calibrate(&calib) {
Some(mut mon) => {
let mut alarm: Option<(usize, struktura::monitor::Leg)> = None;
for (i, row) in stream[calib_n..].iter().enumerate() {
if let Some(leg) = mon.push(row) {
alarm = Some((calib_n + i, leg));
break;
}
}
println!();
println!(" \x1b[1mIMS bearing run-to-failure\x1b[0m ({} recordings, 10 min apart)", n);
println!(" Calibration: recordings 0..{} (healthy)", calib_n);
match alarm {
Some((idx, leg)) => {
let lead_recs = n - 1 - idx;
let lead_hours = lead_recs as f64 * 10.0 / 60.0;
println!(" ALARM at recording {} via {:?}", idx, leg);
println!(" Failure at recording {} (test termination)", n - 1);
println!(" alarm {} recordings ({:.1} h) before the test ended", lead_recs, lead_hours);
println!(" (a plain RMS amplitude threshold on this bearing trips earlier)");
}
None => println!(" \x1b[31mNo alarm raised, detection failed\x1b[0m"),
}
}
None => println!(" IMS calibration failed (insufficient data)"),
}
} else {
println!(" IMS raw data not found (need data/ims/extracted/2nd_test/2nd_test), skipped");
}
if stream.len() >= 900 {
use struktura::prognosis::time_to_threshold;
let n = stream.len();
let b1: Vec<f64> = stream.iter().map(|r| r[0]).collect();
let baseline: f64 = b1[..300].iter().sum::<f64>() / 300.0;
let failure_level = 4.0 * baseline;
let actual_fail = b1[..n - 2]
.iter()
.position(|&v| v > failure_level)
.unwrap_or(n - 3);
println!();
println!(" \x1b[1mPROGNOSIS VALIDATION\x1b[0m (bearing-1 RMS, threshold {:.3} = 4x baseline,",
failure_level);
println!(" actual crossing at recording {})", actual_fail);
println!(" | Checkpoint | Rec | Prediction | Actual | Error |");
println!(" |------------|-----|--------------------------|--------|---------|");
for (label, frac) in [("50%", 0.50), ("75%", 0.75), ("90%", 0.90)] {
let at = (n as f64 * frac) as usize;
let hist = &b1[..at];
match time_to_threshold(hist, 100, failure_level) {
Some(e) => {
let pred = at as f64 + e.eta;
let err = pred - actual_fail as f64;
println!(
" | {:<10} | {} | rec {:.0} (1σ {:.0}..{:.0}) | {} | {:+.0} recs |",
label, at, pred, at as f64 + e.eta_low, at as f64 + e.eta_high,
actual_fail, err
);
}
None => println!(
" | {:<10} | {} | no resolvable trend | {} | -- |",
label, at, actual_fail
),
}
}
println!(" ('no resolvable trend' early in life is the correct output:");
println!(" degradation had not begun, a number there would be fiction.)");
}
let pre: Vec<f64> = include_str!("../../data/voyager1_helio_pre.csv")
.lines().filter_map(|l| l.trim().parse().ok()).collect();
let post: Vec<f64> = include_str!("../../data/voyager1_helio_post.csv")
.lines().filter_map(|l| l.trim().parse().ok()).collect();
let calib_n = pre.len() / 2;
println!();
println!(" \x1b[1mVoyager 1 magnetometer, heliopause crossing\x1b[0m");
println!(" Calibration: first {} pre-crossing samples (2011-2012 cruise);", calib_n);
println!(" stream: remaining {} pre + {} post-crossing samples",
pre.len() - calib_n, post.len());
match HybridMonitor::calibrate(&[pre[..calib_n].to_vec()]) {
Some(mut mon) => {
use struktura::monitor::Leg;
mon.set_leg_enabled(Leg::LevelShift, false);
mon.set_leg_enabled(Leg::ResidualCusum, false);
println!(" Legs: residual + repeated + DFA (level/CUSUM disabled: trending channel)");
let mut alarm_at: Option<(usize, struktura::monitor::Leg)> = None;
let total: Vec<f64> = pre[calib_n..].iter().chain(post.iter()).cloned().collect();
let pre_rest = pre.len() - calib_n;
for (i, &v) in total.iter().enumerate() {
if let Some(leg) = mon.push(&[v]) {
alarm_at = Some((i, leg));
break;
}
}
match alarm_at {
Some((i, leg)) => {
if i < pre_rest {
println!(" \x1b[31mALARM at sample {} PRE-crossing via {:?}, false alarm\x1b[0m", i, leg);
} else {
println!(" Quiet through {} held-out pre-crossing samples;", pre_rest);
println!(" \x1b[32mALARM {} samples after the heliopause crossing via {:?}\x1b[0m",
i - pre_rest, leg);
println!(" (48s cadence: {:.1} hours into interstellar space)",
(i - pre_rest) as f64 * 48.0 / 3600.0);
}
}
None => println!(" No alarm across the crossing"),
}
}
None => println!(" Voyager calibration failed"),
}
println!();
}
fn cmd_generate_hybrid(args: &[String]) {
use struktura::codegen::generate_hybrid_c;
use struktura::monitor::HybridMonitor;
use struktura::telemetry_bench::synth_spacecraft;
let mut out = "hybrid_monitor.c".to_string();
let mut calib_csv: Option<String> = None;
let mut i = 2;
while i < args.len() {
match args[i].as_str() {
"-o" | "--output" if i + 1 < args.len() => { out = args[i + 1].clone(); i += 2; }
"--calib" if i + 1 < args.len() => { calib_csv = Some(args[i + 1].clone()); i += 2; }
_ => { i += 1; }
}
}
let calib: Vec<Vec<f64>> = match calib_csv {
Some(path) => {
let content = std::fs::read_to_string(&path).unwrap_or_else(|e| {
eprintln!("cannot read {}: {}", path, e);
process::exit(1);
});
let rows: Vec<Vec<f64>> = content
.lines()
.filter_map(|l| {
let vals: Vec<f64> =
l.split(',').filter_map(|s| s.trim().parse().ok()).collect();
if vals.is_empty() { None } else { Some(vals) }
})
.collect();
if rows.is_empty() {
eprintln!("no numeric rows in {}", path);
process::exit(1);
}
let nch = rows[0].len();
(0..nch)
.map(|ch| rows.iter().map(|r| r.get(ch).copied().unwrap_or(0.0)).collect())
.collect()
}
None => synth_spacecraft(2048, 424242),
};
let mon = HybridMonitor::calibrate(&calib).unwrap_or_else(|| {
eprintln!("calibration failed (need >= 192 samples per channel, equal lengths)");
process::exit(1);
});
let code = generate_hybrid_c(&mon.export());
std::fs::write(&out, &code).unwrap_or_else(|e| {
eprintln!("cannot write {}: {}", out, e);
process::exit(1);
});
println!("Generated {} ({} channels, {} bytes)", out, mon.channels(), code.len());
println!("Compile: gcc -std=c99 -Wall -Werror -O2 -DHYBRID_STANDALONE_TEST -o hybrid {} -lm", out);
println!("Self-test: ./hybrid (expects: SELFTEST PASS)");
}
fn cmd_mission() {
use struktura::autopilot::{AutoPilot, Event};
use struktura::monitor::HybridMonitor;
use struktura::telemetry_bench::synth_spacecraft;
println!();
println!(" \x1b[1mAUTONOMOUS MISSION GAUNTLET\x1b[0m: no human in the loop");
println!(" 24,000 samples. Scripted events the autopilot must survive alone:");
println!(" t= 4,000 temp sensor freezes (dead sensor)");
println!(" t=10,000 PERMANENT regime change (+0.8 sigma, all channels)");
println!(" t=18,000 drift fault on SOC, in the NEW regime");
println!(" ================================================================");
println!();
let n = 24_000usize;
let calib = synth_spacecraft(2048, 31_337 + 100);
let stream = synth_spacecraft(n, 31_337 + 200);
let ch_sd: Vec<f64> = (0..6)
.map(|ch| {
let c = &calib[ch];
let m = c.iter().sum::<f64>() / c.len() as f64;
(c.iter().map(|x| (x - m) * (x - m)).sum::<f64>() / c.len() as f64).sqrt()
})
.collect();
let mon = HybridMonitor::calibrate(&calib).expect("calibration");
let mut ap = AutoPilot::new(mon);
let valid = [true; 6];
let mut sample = [0.0f64; 6];
let names = ["soc", "bus_voltage", "temp", "wheel", "pointing", "payload_current"];
for t in 0..n {
for ch in 0..6 {
let mut v = stream[ch][t];
if ch == 2 && t >= 4000 { v = stream[2][4000]; }
if t >= 10_000 { v += 0.8 * ch_sd[ch]; }
if ch == 0 && t >= 18_000 { v += (t - 18_000) as f64 * 0.0005 * ch_sd[0]; }
sample[ch] = v;
}
for ev in ap.push(&sample, &valid) {
match ev {
Event::Alarm { tick, report, class } => println!(
" t={:>6} \x1b[33mALARM\x1b[0m {:?} on {} (class: {})",
tick, report.leg, names[report.channel], class
),
Event::Quarantined { tick, channel } => println!(
" t={:>6} \x1b[31mQUARANTINE\x1b[0m {} declared dead -> virtual mode",
tick, names[channel]
),
Event::AdaptationStarted { tick } => println!(
" t={:>6} \x1b[36mADAPTING\x1b[0m level shift: collecting new-regime window",
tick
),
Event::Recalibrated { tick } => println!(
" t={:>6} \x1b[32mRECALIBRATED\x1b[0m guard passed -> new normal accepted",
tick
),
Event::RolledBack { tick, .. } => println!(
" t={:>6} \x1b[31mROLLBACK\x1b[0m guard alarmed -> adaptation rejected",
tick
),
}
}
}
println!();
println!(" Mission complete. Every decision above was made autonomously.");
println!(" (Reproduce exactly: deterministic seeds. Assertions in autopilot tests.)");
println!();
}
fn cmd_redblue(args: &[String]) {
use struktura::redblue::run;
let mut rounds = 6usize;
let mut probes = 120u64;
let mut mutations = 10usize;
let mut i = 2;
while i < args.len() {
match args[i].as_str() {
"--rounds" if i + 1 < args.len() => { rounds = args[i + 1].parse().unwrap_or(6); i += 2; }
"--probes" if i + 1 < args.len() => { probes = args[i + 1].parse().unwrap_or(120); i += 2; }
"--mutations" if i + 1 < args.len() => { mutations = args[i + 1].parse().unwrap_or(10); i += 2; }
_ => { i += 1; }
}
}
println!();
println!(" \x1b[1mRED/BLUE ADVERSARIAL SELF-IMPROVEMENT\x1b[0m");
println!(" RED probes the continuous fault space for misses; BLUE evolves the");
println!(" detection policy against the accumulated miss corpus under a");
println!(" ZERO-clean-alarm law. {} rounds x {} probes x {} mutations.", rounds, probes, mutations);
println!(" ================================================================");
println!();
println!(" | Round | RED coverage | New misses | Corpus cov. after BLUE | Evolved? |");
println!(" |-------|--------------|------------|------------------------|----------|");
let (final_config, reports) = run(rounds, probes, mutations, 12, |r| {
let cov_color = if r.red_coverage >= 0.9 { "\x1b[32m" }
else if r.red_coverage >= 0.7 { "\x1b[33m" } else { "\x1b[31m" };
println!(
" | {:>5} | {}{:>7.1}%\x1b[0m | {:>10} | {:>17.1}% | {:>8} |",
r.round + 1, cov_color, r.red_coverage * 100.0, r.new_misses,
r.corpus_coverage_after * 100.0,
if r.improved { "\x1b[32mYES\x1b[0m" } else { "no" }
);
});
println!();
if let (Some(first), Some(last)) = (reports.first(), reports.last()) {
println!(
" RED coverage: {:.1}% (round 1) -> {:.1}% (round {})",
first.red_coverage * 100.0, last.red_coverage * 100.0, reports.len()
);
}
println!(" Evolved config: res_span={} dfa_persist={} roll_persist={} cusum_k={:.2} horizon={:.0}",
final_config.res_span, final_config.dfa_persist, final_config.roll_persist,
final_config.cusum_k, final_config.design_horizon);
println!(" Zero-clean-alarm law verified on every accepted mutation.");
println!();
}
fn cmd_evolve(args: &[String]) {
use struktura::redblue::evolve;
let mut generations = 10usize;
let mut probes = 100u64;
let mut candidates = 8usize;
let mut i = 2;
while i < args.len() {
match args[i].as_str() {
"--generations" if i + 1 < args.len() => { generations = args[i + 1].parse().unwrap_or(10); i += 2; }
"--probes" if i + 1 < args.len() => { probes = args[i + 1].parse().unwrap_or(100); i += 2; }
"--candidates" if i + 1 < args.len() => { candidates = args[i + 1].parse().unwrap_or(8); i += 2; }
_ => { i += 1; }
}
}
println!();
println!(" \x1b[1mGENERATIONAL EVOLUTION\x1b[0m: parameter mutation + LEG SYNTHESIS");
println!(" BLUE may now compose NEW detector legs from a grammar (source x");
println!(" statistic x window x persistence), each Gumbel-calibrated, accepted");
println!(" only with zero clean alarms. {} generations x {} probes x {} candidates.",
generations, probes, candidates);
println!(" ================================================================");
println!();
println!(" | Gen | RED coverage | New misses | Corpus cov. | Legs | Config (persist d/r, span) |");
println!(" |-----|--------------|------------|-------------|------|-----------------------------|");
let (final_config, final_legs, reports) =
evolve(generations, probes, candidates, 12, 6, |r| {
let c = if r.red_coverage >= 0.9 { "\x1b[32m" }
else if r.red_coverage >= 0.75 { "\x1b[33m" } else { "\x1b[31m" };
println!(
" | {:>3} | {}{:>7.1}%\x1b[0m | {:>10} | {:>9.1}% | {:>4} | d={} r={} span={} |",
r.generation + 1, c, r.red_coverage * 100.0, r.new_misses,
r.corpus_coverage * 100.0, r.legs.len(),
r.config.dfa_persist, r.config.roll_persist, r.config.res_span
);
});
println!();
if let (Some(first), Some(last)) = (reports.first(), reports.last()) {
println!(" RED frontier: {:.1}% (gen 1) -> {:.1}% (gen {})",
first.red_coverage * 100.0, last.red_coverage * 100.0, reports.len());
}
println!(" Synthesized legs: {}", final_legs.len());
for g in &final_legs {
let src = ["raw", "ar1-residual", "first-diff"][g.source as usize % 3];
let st = ["mean", "std", "max-abs", "range", "slope"][g.statistic as usize % 5];
println!(" {} -> {} over {} samples, persist {}",
src, st, [8, 16, 32, 64][g.window as usize % 4],
[1, 2, 4, 8][g.persist as usize % 4]);
}
println!(" Final config: res_span={} dfa_persist={} roll_persist={} cusum_k={:.2} horizon={:.0}",
final_config.res_span, final_config.dfa_persist, final_config.roll_persist,
final_config.cusum_k, final_config.design_horizon);
println!(" Zero-clean-alarm law enforced on every acceptance (12 clean seeds).");
println!();
}
#[allow(dead_code)]
fn read_csv_raw(path: &str) -> (Vec<String>, Vec<Vec<String>>) {
let content = fs::read_to_string(path).unwrap_or_else(|e| {
eprintln!("Error reading {}: {}", path, e);
process::exit(1);
});
let mut lines = content.lines();
let header: Vec<String> = match lines.next() {
Some(h) => h.split(',').map(|s| s.trim().to_string()).collect(),
None => return (vec![], vec![]),
};
let ncols = header.len();
let mut cols: Vec<Vec<String>> = (0..ncols).map(|_| Vec::new()).collect();
for line in lines {
let fields: Vec<&str> = line.split(',').collect();
for i in 0..ncols {
cols[i].push(fields.get(i).map(|s| s.trim().to_string()).unwrap_or_default());
}
}
(header, cols)
}
fn timestamp_row(timestamps: &[String], target: &str) -> usize {
for (i, t) in timestamps.iter().enumerate() {
if t.as_str() >= target {
return i;
}
}
timestamps.len().saturating_sub(1)
}
struct RealLabel {
id: String,
channel: String,
start_time: String,
end_time: String,
}
fn parse_labels_csv(content: &str) -> Vec<RealLabel> {
let mut lines = content.lines();
let header: Vec<String> = match lines.next() {
Some(h) => h.split(',').map(|s| s.trim().to_lowercase()).collect(),
None => return vec![],
};
let idx = |name: &str| header.iter().position(|h| h == name);
let (i_id, i_ch, i_start, i_end) = match (idx("id"), idx("channel"), idx("starttime"), idx("endtime")) {
(Some(a), Some(b), Some(c), Some(d)) => (a, b, c, d),
_ => {
eprintln!("labels CSV must have header ID,Channel,StartTime,EndTime");
process::exit(1);
}
};
let mut out = Vec::new();
for line in lines {
let fields: Vec<&str> = line.split(',').collect();
let get = |i: usize| fields.get(i).map(|s| s.trim().to_string()).unwrap_or_default();
if fields.len() <= i_id.max(i_ch).max(i_start).max(i_end) {
continue;
}
out.push(RealLabel {
id: get(i_id),
channel: get(i_ch),
start_time: get(i_start),
end_time: get(i_end),
});
}
out
}
fn cmd_evolve_real(args: &[String]) {
use struktura::evolve_real::{evolve_real, EvolveRealConfig, RealAnomaly};
let mut train_path = None;
let mut test_path = None;
let mut labels_path = None;
let mut channel_name = None;
let mut rounds = 8usize;
let mut mutations = 8usize;
let mut i = 2;
while i < args.len() {
match args[i].as_str() {
"--train" if i + 1 < args.len() => { train_path = Some(args[i + 1].clone()); i += 2; }
"--test" if i + 1 < args.len() => { test_path = Some(args[i + 1].clone()); i += 2; }
"--labels" if i + 1 < args.len() => { labels_path = Some(args[i + 1].clone()); i += 2; }
"--channel" if i + 1 < args.len() => { channel_name = Some(args[i + 1].clone()); i += 2; }
"--rounds" if i + 1 < args.len() => { rounds = args[i + 1].parse().unwrap_or(8); i += 2; }
"--mutations" if i + 1 < args.len() => { mutations = args[i + 1].parse().unwrap_or(8); i += 2; }
_ => { i += 1; }
}
}
let (_train_path, test_path, labels_path) = match (train_path, test_path, labels_path) {
(Some(a), Some(b), Some(c)) => (a, b, c),
_ => {
eprintln!("usage: struktura evolve-real --train <csv> --test <csv> --labels <csv> [--rounds N] [--mutations N]");
process::exit(1);
}
};
eprintln!("evolve-real: loading timestamps from test file...");
let test_content = fs::read_to_string(&test_path).unwrap_or_else(|e| {
eprintln!("Error reading {}: {}", test_path, e);
process::exit(1);
});
let (test_header, test_cols, _) = parse_csv_columns(&test_content);
let (_, test_raw_cols) = {
let mut lines = test_content.lines();
let hdr: Vec<String> = lines.next().unwrap_or("").split(',').map(|s| s.trim().to_string()).collect();
let ncols = hdr.len();
let mut cols: Vec<Vec<String>> = (0..ncols).map(|_| Vec::new()).collect();
for line in lines {
let fields: Vec<&str> = line.split(',').collect();
for (i, col) in cols.iter_mut().enumerate() {
col.push(fields.get(i).unwrap_or(&"").trim().to_string());
}
}
(hdr, cols)
};
eprintln!("evolve-real: {} rows x {} columns loaded", test_cols.first().map(|c| c.len()).unwrap_or(0), test_header.len());
let train_header = test_header.clone();
let max_calib = 10000usize;
let train_cols: Vec<Vec<f64>> = test_cols.iter().map(|c| {
let start = c.len().saturating_sub(max_calib);
c[start..].to_vec()
}).collect();
let shared: Vec<String> = train_header
.iter()
.filter(|h| test_header.contains(h))
.cloned()
.collect();
if shared.is_empty() {
eprintln!("train/test CSVs share no common column names");
process::exit(1);
}
let calib: Vec<Vec<f64>> = shared
.iter()
.map(|h| train_cols[train_header.iter().position(|x| x == h).unwrap()].clone())
.collect();
let test: Vec<Vec<f64>> = shared
.iter()
.map(|h| test_cols[test_header.iter().position(|x| x == h).unwrap()].clone())
.collect();
let ts_idx = test_header
.iter()
.position(|h| h.to_lowercase().contains("time"))
.unwrap_or(0);
let timestamps = &test_raw_cols[ts_idx];
let labels_content = fs::read_to_string(&labels_path).unwrap_or_else(|e| {
eprintln!("Error reading {}: {}", labels_path, e);
process::exit(1);
});
let labels = parse_labels_csv(&labels_content);
let effective_channel = channel_name.clone().or_else(|| {
if shared.len() == 1 || (shared.len() <= 2 && shared.iter().any(|h| h == "value")) {
None
} else {
Some(String::new()) }
});
let anomalies: Vec<RealAnomaly> = labels
.iter()
.filter(|l| {
match &effective_channel {
Some(ch) if !ch.is_empty() => &l.channel == ch,
Some(_) => shared.iter().any(|h| h == &l.channel), None => channel_name.as_ref().map_or(true, |c| &l.channel == c), }
})
.filter(|l| match &channel_name {
Some(ch) => &l.channel == ch,
None => true,
})
.map(|l| RealAnomaly {
id: l.id.clone(),
start_row: timestamp_row(timestamps, &l.start_time),
end_row: timestamp_row(timestamps, &l.end_time),
})
.collect();
println!();
println!(" \x1b[1mEVOLVE AGAINST REAL LABELED DATA\x1b[0m");
println!(" {} channels matched, {} labeled anomalies loaded ({} rounds x {} mutations).",
shared.len(), anomalies.len(), rounds, mutations);
println!(" ================================================================");
println!();
println!(" | Round | Detected/Total | Coverage | Clean alarms | Improved |");
println!(" |-------|----------------|----------|--------------|----------|");
let cfg = EvolveRealConfig { calib, test, anomalies, rounds, blue_mutations: mutations };
let (final_config, final_legs, reports) = evolve_real(cfg, |r| {
println!(
" | {:>5} | {:>6}/{:<6} | {:>6.1}% | {:>12} | {:>8} |",
r.round + 1, r.detected, r.total, r.coverage * 100.0,
r.clean_false_alarms, if r.improved { "yes" } else { "no" }
);
});
println!();
if let Some(last) = reports.last() {
println!(" Final coverage: {:.1}% ({}/{})", last.coverage * 100.0, last.detected, last.total);
}
println!(" Synthesized legs: {}", final_legs.len());
println!(" {{\"res_span\":{},\"dfa_persist\":{},\"roll_persist\":{},\"cusum_k\":{},\"design_horizon\":{}}}",
final_config.res_span, final_config.dfa_persist, final_config.roll_persist,
final_config.cusum_k, final_config.design_horizon);
println!();
}
fn cmd_smap(args: &[String]) {
use struktura::monitor::HybridMonitor;
let mut ar_mode: Option<usize> = None;
let mut with_dfa = false;
let mut verbose = false;
let mut class_stats: std::collections::BTreeMap<String, (usize, usize)> =
std::collections::BTreeMap::new();
let mut z_min = 2.5f64;
let mut p_prune = 0.13f64;
let mut buffer = 50usize;
let mut i = 2;
while i < args.len() {
match args[i].as_str() {
"--ar" => {
ar_mode = Some(args.get(i + 1).and_then(|s| s.parse().ok()).unwrap_or(25));
i += 2;
}
"--zmin" if i + 1 < args.len() => { z_min = args[i + 1].parse().unwrap_or(2.5); i += 2; }
"--prune" if i + 1 < args.len() => { p_prune = args[i + 1].parse().unwrap_or(0.13); i += 2; }
"--buffer" if i + 1 < args.len() => { buffer = args[i + 1].parse().unwrap_or(50); i += 2; }
"--dfa" => { with_dfa = true; i += 1; }
"--verbose" | "-v" => { verbose = true; i += 1; }
_ => { i += 1; }
}
}
println!();
println!(" \x1b[1mNASA SMAP + MSL/CURIOSITY ANOMALY BENCHMARK\x1b[0m");
println!(" Real spacecraft telemetry, JPL-labeled anomalies (telemanom dataset).");
println!(" Protocol: calibrate on the nominal train split, stream the test");
println!(" split; a labeled sequence is DETECTED if any alarm lands inside it;");
println!(" alarms outside every labeled sequence count as false positives.");
println!(" ================================================================");
println!();
let labels_path = "data/smap_msl/labeled_anomalies.csv";
let labels = match std::fs::read_to_string(labels_path) {
Ok(s) => s,
Err(e) => { eprintln!("cannot read {}: {}", labels_path, e); process::exit(1); }
};
let mut per_sc: Vec<(String, usize, usize, usize, usize, usize)> = vec![
("SMAP".to_string(), 0, 0, 0, 0, 0),
("MSL".to_string(), 0, 0, 0, 0, 0),
];
const REFRACTORY: usize = 100;
for line in labels.lines().skip(1) {
let mut parts = line.splitn(2, ',');
let chan = match parts.next() { Some(c) => c.trim().to_string(), None => continue };
let rest = match parts.next() { Some(r) => r, None => continue };
let mut parts2 = rest.splitn(2, ',');
let spacecraft = match parts2.next() { Some(s) => s.trim().to_string(), None => continue };
let rest2 = match parts2.next() { Some(r) => r, None => continue };
let seq_str: String = {
let start = rest2.find('[').unwrap_or(0);
let end = rest2.rfind(']').map(|e| e + 1).unwrap_or(rest2.len());
rest2[start..end].to_string()
};
let mut sequences: Vec<(usize, usize)> = Vec::new();
{
let nums: Vec<usize> = seq_str
.split(|c: char| !c.is_ascii_digit())
.filter(|s| !s.is_empty())
.filter_map(|s| s.parse().ok())
.collect();
for pair in nums.chunks(2) {
if pair.len() == 2 {
sequences.push((pair[0], pair[1]));
}
}
}
if sequences.is_empty() { continue; }
let sc_idx = if spacecraft == "SMAP" { 0 } else { 1 };
let train_path = format!("data/smap_msl/train_csv/{}.csv", chan);
let test_path = format!("data/smap_msl/test_csv/{}.csv", chan);
let (train, test) = match (std::fs::read_to_string(&train_path), std::fs::read_to_string(&test_path)) {
(Ok(a), Ok(b)) => (a, b),
_ => { per_sc[sc_idx].5 += 1; per_sc[sc_idx].3 += sequences.len(); continue; }
};
let train_v: Vec<f64> = train.lines().filter_map(|l| l.trim().parse().ok()).collect();
let test_v: Vec<f64> = test.lines().filter_map(|l| l.trim().parse().ok()).collect();
if train_v.len() < 300 || test_v.is_empty() {
per_sc[sc_idx].5 += 1;
per_sc[sc_idx].3 += sequences.len();
continue;
}
if let Some(p) = ar_mode {
let mut preds = struktura::smap_eval::detect_channel_tuned(
&train_v, &test_v, p, z_min, p_prune, buffer,
);
if with_dfa {
let dfa_preds = struktura::smap_eval::dfa_sequences(
&train_v, &test_v, 96, z_min, buffer,
);
preds.extend(dfa_preds);
}
let class = line.rsplit(',').nth(1).unwrap_or("?").trim().to_string();
let mut hit = vec![false; sequences.len()];
let mut fp = 0usize;
for &(ps, pe) in &preds {
let mut overlaps = false;
for (si, &(s, e)) in sequences.iter().enumerate() {
if ps <= e && pe >= s {
hit[si] = true;
overlaps = true;
}
}
if !overlaps {
fp += 1;
}
}
for (si, &h) in hit.iter().enumerate() {
let entry = class_stats.entry(class.clone()).or_insert((0usize, 0usize));
entry.1 += 1;
if h {
entry.0 += 1;
}
let _ = si;
}
let tp = hit.iter().filter(|&&h| h).count();
let missed = sequences.len() - tp;
if verbose && missed > 0 {
eprintln!(" MISS {}: {} of {} seqs missed, {} preds, {} FP",
chan, missed, sequences.len(), preds.len(), fp);
}
per_sc[sc_idx].1 += tp;
per_sc[sc_idx].2 += fp;
per_sc[sc_idx].3 += sequences.len();
per_sc[sc_idx].4 += 1;
continue;
}
let trending = {
let n = train_v.len() - 1;
let x = &train_v[..n];
let y = &train_v[1..];
let mx: f64 = x.iter().sum::<f64>() / n as f64;
let my: f64 = y.iter().sum::<f64>() / n as f64;
let mut cov = 0.0;
let mut var = 0.0;
for i in 0..n {
cov += (x[i] - mx) * (y[i] - my);
var += (x[i] - mx) * (x[i] - mx);
}
var < 1e-12 || cov / var > 0.98
};
let mut mon = match HybridMonitor::calibrate(&[train_v]) {
Some(m) => m,
None => { per_sc[sc_idx].5 += 1; per_sc[sc_idx].3 += sequences.len(); continue; }
};
if trending {
use struktura::monitor::Leg;
mon.set_leg_enabled(Leg::LevelShift, false);
mon.set_leg_enabled(Leg::ResidualCusum, false);
}
use struktura::autopilot::{AutoPilot, Event};
use struktura::monitor::Leg;
let mut ap = AutoPilot::new(mon);
let mut alarms: Vec<usize> = Vec::new();
let valid = [true];
for (t, &v) in test_v.iter().enumerate() {
for ev in ap.push(&[v], &valid) {
let is_fault = match &ev {
Event::Alarm { report, class, .. } => {
report.leg != Leg::LevelShift || *class == "drift_confirmed"
}
Event::RolledBack { .. } => true,
_ => false,
};
if is_fault {
if alarms.last().map_or(true, |&last| t >= last + REFRACTORY) {
alarms.push(t);
}
}
}
}
let mut hit = vec![false; sequences.len()];
let mut fp = 0usize;
for &a in &alarms {
let mut inside = false;
for (i, &(s, e)) in sequences.iter().enumerate() {
if a >= s.saturating_sub(0) && a <= e {
hit[i] = true;
inside = true;
}
}
if !inside {
fp += 1;
}
}
let tp = hit.iter().filter(|&&h| h).count();
per_sc[sc_idx].1 += tp;
per_sc[sc_idx].2 += fp;
per_sc[sc_idx].3 += sequences.len();
per_sc[sc_idx].4 += 1;
}
println!(" | Spacecraft | Channels | Sequences | Detected | FP alarms | Precision | Recall |");
println!(" |------------|----------|-----------|----------|-----------|-----------|--------|");
let mut tot = (0usize, 0usize, 0usize);
for (name, tp, fp, seqs, chans, skipped) in &per_sc {
let prec = if tp + fp > 0 { *tp as f64 / (tp + fp) as f64 } else { 0.0 };
let rec = if *seqs > 0 { *tp as f64 / *seqs as f64 } else { 0.0 };
println!(
" | {:<10} | {:>4} ({}sk) | {:>9} | {:>8} | {:>9} | {:>8.1}% | {:>5.1}% |",
name, chans, skipped, seqs, tp, fp, prec * 100.0, rec * 100.0
);
tot.0 += tp; tot.1 += fp; tot.2 += seqs;
}
let prec = if tot.0 + tot.1 > 0 { tot.0 as f64 / (tot.0 + tot.1) as f64 } else { 0.0 };
let rec = if tot.2 > 0 { tot.0 as f64 / tot.2 as f64 } else { 0.0 };
let f1 = if prec + rec > 0.0 { 2.0 * prec * rec / (prec + rec) } else { 0.0 };
println!(" |------------|----------|-----------|----------|-----------|-----------|--------|");
println!(
" | {:<10} | | {:>9} | {:>8} | {:>9} | {:>8.1}% | {:>5.1}% |",
"TOTAL", tot.2, tot.0, tot.1, prec * 100.0, rec * 100.0
);
println!();
println!(" Overall F1: {:.3} (JPL telemanom LSTM, same data: P=87.5% R=80.0%)", f1);
if !class_stats.is_empty() {
println!(" Recall by anomaly class:");
for (class, (hits, total)) in &class_stats {
println!(" {:<12} {:>3}/{:<3} = {:>5.1}%",
class, hits, total, *hits as f64 / *total as f64 * 100.0);
}
}
println!(" Self-calibrated, no training, no GPU, ~4us/sample. Note: the values are");
println!(" pre-scaled to (-1,1) by JPL and many channels saturate, which auto-disables");
println!(" the repeated-value leg on those channels.");
println!();
}
fn cmd_nasa() {
use struktura::smap_eval::detect_channel_tuned;
let train: Vec<f64> = include_str!("../../data/smap_t1_train.csv")
.lines().filter_map(|l| l.trim().parse().ok()).collect();
let test: Vec<f64> = include_str!("../../data/smap_t1_test.csv")
.lines().filter_map(|l| l.trim().parse().ok()).collect();
println!();
println!(" \x1b[1mNASA SMAP SATELLITE, CHANNEL T-1 (thermal telemetry)\x1b[0m");
println!(" Real data from NASA's Soil Moisture Active Passive satellite.");
println!(" Two labeled anomalies: a point fault + a contextual fault.");
println!(" Method: ridge AR(auto) + JPL's telemanom residual math.");
println!(" ================================================================");
println!();
let seqs = detect_channel_tuned(&train, &test, 0, 0.5, 0.02, 150);
use struktura::monitor::HybridMonitor;
let mut mon = HybridMonitor::calibrate(std::slice::from_ref(&train)).expect("calibration");
mon.set_leg_enabled(struktura::monitor::Leg::LevelShift, false);
mon.set_leg_enabled(struktura::monitor::Leg::ResidualCusum, false);
let mut stream_alarms: Vec<(usize, struktura::monitor::Leg)> = Vec::new();
for (t, &v) in test.iter().enumerate() {
if let Some(leg) = mon.push(&[v]) {
stream_alarms.push((t, leg));
mon.reset();
}
}
let labeled = [(2399usize, 3898usize), (6550, 6585)];
let classes = ["point", "contextual"];
for (i, &(ls, le)) in labeled.iter().enumerate() {
let found = seqs.iter().any(|&(s, e)| s <= le && e >= ls);
let mark = if found { "\x1b[32m✓ DETECTED\x1b[0m" } else { "\x1b[31m✗ MISSED\x1b[0m" };
println!(" Anomaly {}: samples {}..{} ({}), {}", i + 1, ls, le, classes[i], mark);
}
let fp: Vec<_> = seqs.iter().filter(|&&(s, e)| {
!labeled.iter().any(|&(ls, le)| s <= le && e >= ls)
}).collect();
println!();
println!(" Detected sequences: {:?}", seqs);
println!(" False positives: {}", fp.len());
println!();
println!(" Train: {} samples (nominal). Test: {} samples.", train.len(), test.len());
println!(" Calibration: closed-form ridge AR, no gradient training, no GPU.");
println!();
println!(" \x1b[1mStreaming monitor\x1b[0m (7 legs, self-calibrated, 4us/sample):");
for (t, leg) in &stream_alarms {
let in_anomaly = labeled.iter().any(|&(s, e)| *t >= s && *t <= e);
let mark = if in_anomaly { "\x1b[32m(TRUE)\x1b[0m" } else { "\x1b[31m(FP)\x1b[0m" };
println!(" t={:>5} {:?} {}", t, leg, mark);
}
let stream_tp = stream_alarms.iter().filter(|(t, _)| {
labeled.iter().any(|&(s, e)| *t >= s && *t <= e)
}).count();
println!(" {} alarms: {} in anomaly windows, {} false positives",
stream_alarms.len(), stream_tp, stream_alarms.len() - stream_tp);
println!();
println!(" Full benchmark (82 channels): struktura smap --ar 0 (needs data/)");
println!();
}
fn cmd_guard(args: &[String]) {
let mut file_path = String::new();
let mut baseline_n = 0usize;
let mut json = false;
let mut watch = false;
let mut watch_ms = 1000u64;
let mut webhook_url = String::new();
let mut quiet = false;
let mut high = false;
let mut i = 2;
while i < args.len() {
match args[i].as_str() {
"--baseline" if i + 1 < args.len() => { baseline_n = args[i + 1].parse().unwrap_or(0); i += 2; }
"--json" => { json = true; i += 1; }
"--quiet-drift" => { quiet = true; i += 1; }
"--sensitivity" if i + 1 < args.len() => {
high = match args[i + 1].as_str() {
"normal" => false,
"high" => true,
other => {
eprintln!("--sensitivity must be normal or high, got {}", other);
process::exit(2);
}
};
i += 2;
}
"--watch" | "-w" => { watch = true; i += 1; }
"--interval" if i + 1 < args.len() => { watch_ms = args[i + 1].parse().unwrap_or(1000); i += 2; }
"--webhook" if i + 1 < args.len() => { webhook_url = args[i + 1].clone(); i += 2; }
"--help" | "-h" => {
println!("struktura guard <file.csv> [--baseline N] [--sensitivity normal|high] [--json] [--watch] [--webhook URL] [--quiet-drift]");
println!(" Monitor any CSV for anomalies. Exit: 0=healthy 1=fault 2=error");
println!(" --sensitivity normal (default): fewest false alarms. high: catches more, alarms more");
println!(" (NAB: 36 -> 49 of 116 windows, 35 -> 48 false alarms; clean slow-wander");
println!(" synthetic streams 0 -> 2-3 of 30)");
println!(" --quiet-drift Clip what the drift leg sees; small effect (NAB: 35 -> 33 false alarms,");
println!(" 36 -> 35 windows); a spike only the drift leg catches is found later or not at all");
println!(" --watch Follow the file (like tail -f), monitor new rows live");
println!(" --interval MS Poll interval for --watch (default 1000ms)");
println!(" --webhook URL POST anomaly alerts to a Slack/Discord/PagerDuty webhook");
process::exit(0);
}
s if !s.starts_with('-') && file_path.is_empty() => { file_path = s.to_string(); i += 1; }
_ => { i += 1; }
}
}
let content = if file_path.is_empty() || file_path == "-" {
use std::io::Read;
let mut buf = String::new();
std::io::stdin().read_to_string(&mut buf).expect("read stdin");
buf
} else {
std::fs::read_to_string(&file_path).unwrap_or_else(|e| {
eprintln!("cannot read {}: {}", file_path, e);
process::exit(2);
})
};
if !webhook_url.is_empty() {
std::env::set_var("STRUKTURA_WEBHOOK", &webhook_url);
}
let cfg = guard_config(quiet, high);
if watch && !file_path.is_empty() && file_path != "-" {
run_guard_watch(&file_path, baseline_n, json, watch_ms, cfg);
}
let exit = run_guard(&content, baseline_n, json, cfg);
process::exit(exit);
}
const HIGH_SENSITIVITY_HORIZON: f64 = 1e5;
fn guard_config(quiet: bool, high: bool) -> struktura::monitor::MonitorConfig {
let mut cfg = struktura::monitor::MonitorConfig { quiet_drift: quiet, ..Default::default() };
if high {
cfg.design_horizon = HIGH_SENSITIVITY_HORIZON;
}
cfg
}
fn run_guard_watch(path: &str, baseline_n: usize, json: bool, poll_ms: u64, cfg: struktura::monitor::MonitorConfig) -> ! {
use struktura::monitor::HybridMonitor;
use struktura::autopilot::AutoPilot;
let initial = std::fs::read_to_string(path).unwrap_or_else(|e| {
eprintln!("cannot read {}: {}", path, e);
process::exit(2);
});
let parsed = parse_multi_csv(&initial);
let (keep, dropped_cols) = parsed.sensor_columns();
let rows: Vec<Vec<f64>> = parsed
.rows
.iter()
.map(|r| keep.iter().map(|&c| r.get(c).copied().unwrap_or(0.0)).collect())
.collect();
if rows.is_empty() { eprintln!("no numeric rows"); process::exit(2); }
let ncols = keep.len();
let calib_n = if baseline_n > 0 { baseline_n.min(rows.len()) } else { rows.len() };
let channels: Vec<Vec<f64>> = (0..ncols)
.map(|ch| rows.iter().map(|r| r.get(ch).copied().unwrap_or(0.0)).collect())
.collect();
let calib: Vec<Vec<f64>> = channels.iter().map(|c| c[..calib_n].to_vec()).collect();
let mon = match HybridMonitor::calibrate_with(&calib, cfg) {
Some(m) => m,
None => { eprintln!("calibration failed (need >= 192 samples)"); process::exit(2); }
};
let mut ap = AutoPilot::new(mon);
let valid: Vec<bool> = vec![true; ncols];
let mut sample = vec![0.0f64; ncols];
let mut last_alarm: Vec<(usize, u8)> = Vec::new();
const ALARM_COOLDOWN: usize = 50;
let mut emit_dedup = |t: usize, ev: &struktura::autopilot::Event, json: bool| {
if let struktura::autopilot::Event::Alarm { report, .. } = ev {
let leg_id = report.leg as u8;
let dup = last_alarm.iter().any(|&(lt, ll)| ll == leg_id && t.saturating_sub(lt) < ALARM_COOLDOWN);
last_alarm.retain(|&(lt, _)| t.saturating_sub(lt) < ALARM_COOLDOWN);
last_alarm.push((t, leg_id));
if dup { return; }
}
emit_event(t, ev, json);
};
for t in calib_n..rows.len() {
for ch in 0..ncols { sample[ch] = rows[t].get(ch).copied().unwrap_or(0.0); }
for ev in ap.push(&sample, &valid) {
emit_dedup(t, &ev, json);
}
}
let mut lines_seen = initial.lines().count();
if !json {
eprintln!("struktura guard --watch: {} ch, calibrated on {} rows, tailing {} (poll {}ms, Ctrl-C to stop)",
ncols, calib_n, path, poll_ms);
index_columns_note(&dropped_cols);
short_calibration_note(calib_n);
calibration_self_check_note(calibration_self_check(&calib, cfg), calib_n);
}
let mut t = rows.len();
loop {
std::thread::sleep(std::time::Duration::from_millis(poll_ms));
let content = match std::fs::read_to_string(path) {
Ok(c) => c,
Err(_) => continue,
};
let new_lines: Vec<&str> = content.lines().skip(lines_seen).collect();
if new_lines.is_empty() { continue; }
lines_seen += new_lines.len();
for line in new_lines {
let vals: Vec<f64> = line.split(',').filter_map(|s| s.trim().parse().ok()).collect();
if !vals.is_empty() {
for (ch, &col) in keep.iter().enumerate() { sample[ch] = vals.get(col).copied().unwrap_or(0.0); }
for ev in ap.push(&sample, &valid) {
emit_dedup(t, &ev, json);
}
t += 1;
}
}
}
}
fn generate_rover_c_wrapper() -> String {
let mut s = String::new();
s.push_str("/* RoverHealthImpl.c, generated by struktura generate --rover\n");
s.push_str(" * Bridges F Prime component calls to the Rust rover_flight staticlib.\n");
s.push_str(" * Link with: libstruktura.a (cargo build --release --lib --crate-type staticlib)\n");
s.push_str(" */\n\n");
s.push_str("#include <stdint.h>\n\n");
s.push_str("/* FFI declarations, match rover_flight.rs exports */\n");
s.push_str("#define N_CH 10\n\n");
s.push_str("typedef struct { uint8_t leg; uint8_t channel; uint32_t tick; } FAlarm;\n\n");
s.push_str("/* Extern: the Rust staticlib provides these */\n");
s.push_str("extern void rover_monitor_init(void);\n");
s.push_str("extern void rover_monitor_calibrate_channel(unsigned ch, double ar_a, double ar_b, double ar_sd, double mean, double roll_max_dev, unsigned max_run, int repeat_enabled);\n");
s.push_str("extern void rover_monitor_set_res_threshold(double thr);\n");
s.push_str("extern int rover_monitor_push(const double sample[N_CH], FAlarm *out);\n");
s.push_str("extern void rover_monitor_reset(void);\n\n");
s.push_str("/* F Prime component implementation */\n");
s.push_str("static int initialized = 0;\n\n");
s.push_str("void RoverHealth_schedIn_handler(void) {\n");
s.push_str(" if (!initialized) {\n");
s.push_str(" rover_monitor_init();\n");
s.push_str(" /* TODO: load calibration from EEPROM or ground-uploaded params */\n");
s.push_str(" rover_monitor_set_res_threshold(15.0);\n");
s.push_str(" initialized = 1;\n");
s.push_str(" }\n\n");
s.push_str(" /* TODO: read current telemetry from input ports into sample[] */\n");
s.push_str(" double sample[N_CH] = {0};\n");
s.push_str(" /* sample[0] = wheelFL_current; ... etc */\n\n");
s.push_str(" FAlarm alarm;\n");
s.push_str(" if (rover_monitor_push(sample, &alarm)) {\n");
s.push_str(" switch (alarm.leg) {\n");
s.push_str(" case 0: /* Residual */\n");
s.push_str(" /* log_WARNING_HI_Drift(channel_name(alarm.channel)); */\n");
s.push_str(" break;\n");
s.push_str(" case 1: /* Stuck */\n");
s.push_str(" /* log_WARNING_HI_SensorStuck(channel_name(alarm.channel)); */\n");
s.push_str(" break;\n");
s.push_str(" case 2: /* LevelShift */\n");
s.push_str(" /* log_WARNING_HI_LevelShift(channel_name(alarm.channel)); */\n");
s.push_str(" break;\n");
s.push_str(" }\n");
s.push_str(" rover_monitor_reset();\n");
s.push_str(" }\n");
s.push_str("}\n");
s
}
fn cmd_stamp(args: &[String]) {
if args.len() < 3 {
eprintln!("Usage: struktura stamp <file.csv>");
eprintln!(" Write a structural fingerprint into the CSV header.");
eprintln!(" guard/check auto-read the stamp as the baseline.");
process::exit(1);
}
let path = &args[2];
let content = std::fs::read_to_string(path).unwrap_or_else(|e| {
eprintln!("cannot read {}: {}", path, e);
process::exit(1);
});
let data = read_csv(path);
if data.len() < 64 {
eprintln!("need >= 64 samples, got {}", data.len());
process::exit(1);
}
let law = analyze(&data);
let stamp = format!(
"# struktura:baseline alpha={:.4} r2={:.4} hurst={:.4} n={} kurtosis={:.2}\n",
law.dfa.alpha, law.dfa.r_squared, law.hurst, law.n, law.kurtosis
);
let mut out = stamp;
for line in content.lines() {
if line.starts_with("# struktura:baseline") { continue; }
out.push_str(line);
out.push('\n');
}
std::fs::write(path, &out).unwrap_or_else(|e| {
eprintln!("cannot write {}: {}", path, e);
process::exit(1);
});
eprintln!("stamped {}: alpha={:.4} r2={:.4} n={}", path, law.dfa.alpha, law.dfa.r_squared, law.n);
}
fn read_stamp(content: &str) -> Option<f64> {
let first = content.lines().next()?.trim_start_matches('\u{feff}');
if !first.starts_with("# struktura:baseline") { return None; }
first.split_whitespace()
.find(|s| s.starts_with("alpha="))
.and_then(|s| s.strip_prefix("alpha="))
.and_then(|s| s.parse().ok())
}
fn cmd_rover() {
use struktura::rover::{RoverSim, RoverFault, ROVER_CHANNELS, ROVER_CHANNEL_NAMES};
use struktura::monitor::HybridMonitor;
use struktura::autopilot::{AutoPilot, Event};
use struktura::monitor::explain_alarm;
println!();
println!(" \x1b[1mROVER HEALTH MONITORING\x1b[0m: autonomous fault handling");
println!(" 10 channels: 4 wheel motors, suspension, battery (V + SOC),");
println!(" thermal (CPU + motors), comms. three scripted faults:");
println!(" t=1500 front-left wheel bearing starts wearing");
println!(" t=2200 rear-right wheel stalls completely");
println!(" t=2600 battery cell degradation");
println!(" ================================================================");
println!();
let mut sim = RoverSim::new(42);
sim.inject(1500, RoverFault::WheelBearing { wheel: 0, severity: 0.9 });
sim.inject(2200, RoverFault::WheelStall { wheel: 3 });
sim.inject(2600, RoverFault::BatteryCell { severity: 0.8 });
let data = sim.run(3000);
let calib: Vec<Vec<f64>> = data.iter().map(|c| c[..1000].to_vec()).collect();
let mon = HybridMonitor::calibrate(&calib).expect("calibration");
let mut ap = AutoPilot::new(mon);
let valid = [true; ROVER_CHANNELS];
let mut sample = [0.0f64; ROVER_CHANNELS];
let mut _last_leg = 255u8;
let mut last_t = 0usize;
let mut alarm_count = 0usize;
let mut q_count = 0usize;
for t in 1000..3000 {
for ch in 0..ROVER_CHANNELS { sample[ch] = data[ch][t]; }
for ev in ap.push(&sample, &valid) {
match &ev {
Event::Alarm { report, .. } => {
let lid = report.leg as u8;
if t - last_t < 200 { continue; }
_last_leg = lid; last_t = t; alarm_count += 1;
let ch_name = ROVER_CHANNEL_NAMES.get(report.channel).unwrap_or(&"?");
println!(" t={:>5} ⚠ {}: {}", t, ch_name, explain_alarm(report));
}
Event::Quarantined { channel, .. } => {
q_count += 1;
let ch_name = ROVER_CHANNEL_NAMES.get(*channel).unwrap_or(&"?");
println!(" t={:>5} ✗ {} dead, virtual readings active", t, ch_name);
}
Event::AdaptationStarted { .. } =>
println!(" t={:>5} ↻ environment change? learning new baseline...", t),
Event::Recalibrated { .. } =>
println!(" t={:>5} ✓ new baseline accepted", t),
Event::RolledBack { .. } =>
println!(" t={:>5} ✗ not environment, fault confirmed", t),
}
}
}
println!();
println!(" {} anomalies, {} channels quarantined. all decisions autonomous.", alarm_count, q_count);
println!(" simulated rover, scripted faults above; not Spirit or other flight data.");
println!();
}
const MIN_RELIABLE_CALIB: usize = 768;
fn short_calibration_note(calib_n: usize) {
if calib_n < MIN_RELIABLE_CALIB {
eprintln!(" note: {} calibration rows is short; use --baseline {} or more if the data allows (fewer raises level-shift false alarms)",
calib_n, MIN_RELIABLE_CALIB);
}
}
fn calibration_self_check(calib: &[Vec<f64>], config: struktura::monitor::MonitorConfig) -> Option<usize> {
use struktura::monitor::HybridMonitor;
let n = calib.first()?.len();
let half = n / 2;
if half < MIN_RELIABLE_CALIB {
return None;
}
let first: Vec<Vec<f64>> = calib.iter().map(|c| c[..half].to_vec()).collect();
let mut mon = HybridMonitor::calibrate_with(&first, config)?;
let mut sample = vec![0.0f64; calib.len()];
for t in half..n {
for (ch, c) in calib.iter().enumerate() {
sample[ch] = c[t];
}
if mon.push(&sample).is_some() {
return Some(t);
}
}
None
}
fn calibration_self_check_note(row: Option<usize>, calib_n: usize) {
if let Some(row) = row {
eprintln!(" warning: the calibration rows (0..{}) may contain a fault: calibrating on their first half, row {} already alarms.", calib_n, row);
eprintln!(" guard treats calibration rows as healthy. Pass --baseline N with a stretch you know is healthy.");
}
}
fn emit_event(t: usize, ev: &struktura::autopilot::Event, json: bool) {
use struktura::autopilot::Event;
use struktura::monitor::explain_alarm;
match ev {
Event::Alarm { report, class, .. } => {
let explanation = explain_alarm(report);
if json {
println!("{{\"event\":\"alarm\",\"t\":{},\"channel\":{},\"class\":\"{}\",\"explanation\":\"{}\"}}",
t, report.channel, class, explanation);
} else {
eprintln!(" row {:>6} ⚠ ch{}: {}", t, report.channel, explanation);
}
}
Event::Quarantined { channel, .. } => {
if json { println!("{{\"event\":\"quarantine\",\"t\":{},\"channel\":{}}}", t, channel); }
else { eprintln!(" row {:>6} ✗ ch{} declared dead, using reconstructed values", t, channel); }
}
Event::AdaptationStarted { .. } => {
if json { println!("{{\"event\":\"adapting\",\"t\":{}}}", t); }
else { eprintln!(" row {:>6} ↻ environment may have changed, learning new baseline...", t); }
}
Event::Recalibrated { .. } => {
if json { println!("{{\"event\":\"recalibrated\",\"t\":{}}}", t); }
else { eprintln!(" row {:>6} ✓ new baseline accepted, this is the new normal", t); }
}
Event::RolledBack { .. } => {
if json { println!("{{\"event\":\"rollback\",\"t\":{}}}", t); }
else { eprintln!(" row {:>6} ✗ not a real environment change, fault confirmed", t); }
}
}
}
fn run_guard(content: &str, baseline_n: usize, json: bool, cfg: struktura::monitor::MonitorConfig) -> i32 {
use struktura::monitor::HybridMonitor;
use struktura::autopilot::{AutoPilot, Event};
let parsed = parse_multi_csv(content);
let (keep, dropped_cols) = parsed.sensor_columns();
let col_names: Vec<String> = if parsed.col_names.is_empty() {
Vec::new()
} else {
keep.iter().map(|&c| parsed.ch_name(c)).collect()
};
let rows: Vec<Vec<f64>> = parsed
.rows
.iter()
.map(|r| keep.iter().map(|&c| r.get(c).copied().unwrap_or(0.0)).collect())
.collect();
if rows.is_empty() {
if json {
println!("{{\"error\":\"no numeric rows\"}}");
} else {
eprintln!("no numeric rows");
}
return 2;
}
let n = rows.len();
let ncols = rows[0].len();
let calib_n = if baseline_n > 0 { baseline_n } else { (n / 3).clamp(192, 100_000) }.min(n);
let channels: Vec<Vec<f64>> = (0..ncols)
.map(|ch| rows.iter().map(|r| r.get(ch).copied().unwrap_or(0.0)).collect())
.collect();
let calib: Vec<Vec<f64>> = channels.iter().map(|c| c[..calib_n].to_vec()).collect();
let mon = match HybridMonitor::calibrate_with(&calib, cfg) {
Some(m) => m,
None => {
if json {
println!("{{\"error\":\"calibration failed\",\"samples\":{},\"min_required\":192}}", calib_n);
} else {
eprintln!("calibration failed (need >= 192 samples, got {})", calib_n);
}
return 2;
}
};
if !json {
if col_names.is_empty() {
eprintln!("struktura guard: {} samples x {} channels, calibrated on {} rows", n, ncols, calib_n);
} else {
eprintln!("struktura guard: {} samples x {} channels ({}), calibrated on {} rows",
n, ncols, col_names.join(", "), calib_n);
}
index_columns_note(&dropped_cols);
short_calibration_note(calib_n);
}
let calib_suspect = calibration_self_check(&calib, cfg);
if !json {
calibration_self_check_note(calib_suspect, calib_n);
}
let mut sample = vec![0.0f64; ncols];
let mut ap = AutoPilot::new(mon);
let ch_name = |idx: usize| -> String {
col_names.get(idx).cloned().unwrap_or_else(|| format!("ch{}", idx))
};
let mut alarm_count = 0usize;
let mut adapt_count = 0usize;
let mut quarantine_count = 0usize;
let valid: Vec<bool> = vec![true; ncols];
let mut last_alarm: Vec<(usize, u8)> = Vec::new(); const ALARM_COOLDOWN: usize = 50;
for t in calib_n..n {
for ch in 0..ncols {
sample[ch] = channels[ch][t];
}
for ev in ap.push(&sample, &valid) {
match &ev {
Event::Alarm { report, class, .. } => {
let leg_id = report.leg as u8;
let dup = last_alarm.iter().any(|&(lt, ll)| ll == leg_id && t - lt < ALARM_COOLDOWN);
last_alarm.retain(|&(lt, _)| t - lt < ALARM_COOLDOWN);
last_alarm.push((t, leg_id));
if dup { continue; }
alarm_count += 1;
let explanation = struktura::monitor::explain_alarm(report);
let name = ch_name(report.channel);
let score_ratio = report.observed / report.threshold.max(1e-12);
if json {
println!("{{\"event\":\"alarm\",\"t\":{},\"channel\":\"{}\",\"class\":\"{}\",\"score_ratio\":{:.2},\"explanation\":\"{}\"}}",
t, name, class, score_ratio, explanation);
} else {
eprintln!(" row {:>6} ⚠ {} ({:.1}x threshold): {}", t, name, score_ratio, explanation);
}
}
Event::Quarantined { channel, .. } => {
quarantine_count += 1;
let name = ch_name(*channel);
if json {
println!("{{\"event\":\"quarantine\",\"t\":{},\"channel\":\"{}\"}}", t, name);
} else {
eprintln!(" row {:>6} ✗ {} declared dead, using reconstructed values", t, name);
}
}
Event::AdaptationStarted { .. } => {
if json {
println!("{{\"event\":\"adapting\",\"t\":{}}}", t);
} else {
eprintln!(" row {:>6} ↻ environment may have changed, learning new baseline...", t);
}
}
Event::Recalibrated { .. } => {
adapt_count += 1;
if json {
println!("{{\"event\":\"recalibrated\",\"t\":{}}}", t);
} else {
eprintln!(" row {:>6} ✓ new baseline accepted, this is the new normal", t);
}
}
Event::RolledBack { .. } => {
if json {
println!("{{\"event\":\"rollback\",\"t\":{}}}", t);
} else {
eprintln!(" row {:>6} ✗ not a real environment change, fault confirmed", t);
}
}
}
}
}
if json {
println!("{{\"summary\":true,\"samples\":{},\"channels\":{},\"baseline\":{},\"alarms\":{},\"adaptations\":{},\"quarantines\":{},\"verdict\":\"{}\",\"calibration_suspect_row\":{}}}",
n, ncols, calib_n, alarm_count, adapt_count, quarantine_count,
if alarm_count == 0 { "healthy" } else { "fault_detected" },
calib_suspect.map_or("null".to_string(), |r| r.to_string()));
} else {
if alarm_count == 0 && calib_suspect.is_some() {
eprintln!(" no anomalies in the {} samples after calibration, but see the calibration warning above", n - calib_n);
} else if alarm_count == 0 {
eprintln!(" HEALTHY, {} samples, no anomalies", n - calib_n);
} else {
eprintln!(" {} faults detected across {} samples ({} adaptations, {} quarantines)",
alarm_count, n - calib_n, adapt_count, quarantine_count);
}
}
if alarm_count > 0 { 1 } else { 0 }
}
fn cmd_when(args: &[String]) {
use struktura::changepoint::find_changepoints;
if args.len() < 3 {
eprintln!("Usage: struktura when <file.csv> [--max N] [--truth a-b,c-d]");
eprintln!(" Find WHEN something changed in your data. --max caps the number reported (default 5).");
process::exit(1);
}
let path = &args[2];
let max = args.iter().position(|a| a == "--max")
.and_then(|i| args.get(i + 1))
.and_then(|s| s.parse().ok())
.unwrap_or(5usize);
let json = args.iter().any(|a| a == "--json");
let data = if path == "-" { read_stdin() } else { read_csv(path) };
if data.len() < struktura::changepoint::MIN_SAMPLES {
eprintln!("need >= {} samples for changepoint detection, got {}", struktura::changepoint::MIN_SAMPLES, data.len());
process::exit(1);
}
let cps = find_changepoints(&data, 128, max);
let truth: Vec<(usize, usize)> = args.iter().position(|a| a == "--truth")
.and_then(|i| args.get(i + 1))
.map(|s| s.split(',').filter_map(|w| {
let (a, b) = w.split_once('-')?;
Some((a.trim().parse().ok()?, b.trim().parse().ok()?))
}).collect())
.unwrap_or_default();
let in_truth = |loc: usize| truth.iter().any(|&(a, b)| loc >= a && loc <= b);
let hits = cps.iter().filter(|cp| in_truth(cp.location)).count();
if json {
print!("[");
for (i, cp) in cps.iter().enumerate() {
if i > 0 { print!(","); }
print!("{{\"at\":{},\"alpha_before\":{:.3},\"alpha_after\":{:.3},\"shift\":{:.3},\"z\":{:.1}",
cp.location, cp.alpha_before, cp.alpha_after, cp.shift, cp.confidence);
if !truth.is_empty() { print!(",\"in_truth\":{}", in_truth(cp.location)); }
print!("}}");
}
println!("]");
return;
}
println!();
if cps.is_empty() {
println!(" no structural changes detected in {} samples.", data.len());
println!(" the signal's correlation structure is consistent throughout.");
} else {
println!(" \x1b[1m{} structural change{} found in {} samples:\x1b[0m",
cps.len(), if cps.len() == 1 { "" } else { "s" }, data.len());
println!();
for (i, cp) in cps.iter().enumerate() {
let before_interp = interpret_alpha(cp.alpha_before);
let after_interp = interpret_alpha(cp.alpha_after);
let direction = if cp.shift > 0.0 { "more correlated" } else { "less correlated" };
println!(" {}. \x1b[33msample {}\x1b[0m, structure shifted {:+.3} (z={:.1})", i + 1, cp.location, cp.shift, cp.confidence);
println!(" before: α={:.3} ({})", cp.alpha_before, before_interp);
println!(" after: α={:.3} ({})", cp.alpha_after, after_interp);
let mark = if truth.is_empty() { "" } else if in_truth(cp.location) { " [in truth window]" } else { " [outside truth]" };
println!(" → the signal became {} after this point{}", direction, mark);
println!();
}
}
if !truth.is_empty() {
let pct = if cps.is_empty() { 0.0 } else { 100.0 * hits as f64 / cps.len() as f64 };
println!(" truth windows: {} of {} changepoints inside ({:.0}%)", hits, cps.len(), pct);
println!();
}
}
fn cmd_prove(args: &[String]) {
if args.len() < 3 {
eprintln!("Usage: struktura prove <file.csv> [--resamples N] [--json]");
process::exit(1);
}
let path = &args[2];
let n_res = args.iter().position(|a| a == "--resamples")
.and_then(|i| args.get(i + 1))
.and_then(|s| s.parse().ok())
.unwrap_or(100usize);
let json = args.iter().any(|a| a == "--json");
let data = if path == "-" { read_stdin() } else { read_csv(path) };
if data.len() < 64 {
eprintln!("need >= 64 samples, got {}", data.len());
process::exit(1);
}
let ci = bootstrap_alpha(&data, n_res);
let proof = prove_structure(&data);
if json {
println!("{{\"alpha\":{:.4},\"ci_low\":{:.4},\"ci_high\":{:.4},\"resamples\":{},\"shuffled_alpha\":{:.4},\"structure_confirmed\":{}}}",
ci.alpha, ci.ci_low, ci.ci_high, ci.n_resamples, proof.shuffled_alpha, proof.structure_confirmed);
return;
}
println!();
println!(" \x1b[1mstruktura prove\x1b[0m {} samples", data.len());
println!(" alpha: {:.3} 95% CI [{:.3}, {:.3}] ({} bootstrap resamples)", ci.alpha, ci.ci_low, ci.ci_high, ci.n_resamples);
println!(" shuffled: {:.3} (ordering destroyed; real structure collapses toward 0.5)", proof.shuffled_alpha);
println!(" verdict: {}", if proof.structure_confirmed { "STRUCTURE CONFIRMED (alpha is a property of the ordering, not the values)" } else { "NOT CONFIRMED (shuffled alpha is as far from 0.5 as the real one)" });
println!();
}
fn interpret_alpha(alpha: f64) -> &'static str {
if alpha < 0.4 { "anti-correlated / mean-reverting" }
else if alpha < 0.55 { "uncorrelated / random" }
else if alpha < 0.75 { "weakly correlated" }
else if alpha < 1.05 { "strongly correlated (1/f)" }
else if alpha < 1.4 { "highly persistent" }
else { "unbounded walk" }
}
fn generate_ros_package() -> String {
let mut s = String::new();
s.push_str("<?xml version=\"1.0\"?>\n");
s.push_str("<package format=\"3\">\n");
s.push_str(" <name>dfa_monitor</name>\n <version>1.0.0</version>\n");
s.push_str(" <description>DFA structural health monitor for ROS 2</description>\n");
s.push_str(" <maintainer email=\"maintainer@example.com\">koscak-labs</maintainer>\n");
s.push_str(" <license>MIT</license>\n\n");
s.push_str(" <buildtool_depend>ament_cmake</buildtool_depend>\n");
s.push_str(" <depend>rclcpp</depend>\n <depend>std_msgs</depend>\n\n");
s.push_str(" <export>\n <build_type>ament_cmake</build_type>\n </export>\n");
s.push_str("</package>\n");
s
}
fn cmd_pipe(args: &[String]) {
use std::io::{self, BufRead};
use struktura::dfa;
let window: usize = args.iter().position(|a| a == "--window")
.and_then(|i| args.get(i + 1))
.and_then(|s| s.parse().ok())
.unwrap_or(256);
let json = args.iter().any(|a| a == "--json");
if args.iter().any(|a| a == "--help" || a == "-h") {
println!("struktura pipe [--window N] [--json]");
println!(" stream DFA from stdin. one number per line. outputs per-window alpha.");
println!();
println!(" examples:");
println!(" curl -s prometheus:9090/metrics | struktura pipe");
println!(" tail -f /var/log/latency.csv | struktura pipe --window 128");
println!(" mqtt sub sensor/temp | struktura pipe --json");
process::exit(0);
}
let stdin = io::stdin();
let mut buf: Vec<f64> = Vec::with_capacity(window);
let mut baseline_alpha: Option<f64> = None;
for line in stdin.lock().lines() {
let line = match line { Ok(l) => l, Err(_) => break };
let val: f64 = match line.trim().parse::<f64>() {
Ok(v) if v.is_finite() => v,
_ => continue,
};
buf.push(val);
if buf.len() > window { buf.remove(0); }
if buf.len() < window { continue; }
let r = dfa(&buf);
let alpha = r.alpha;
let r2 = r.r_squared;
if baseline_alpha.is_none() {
baseline_alpha = Some(alpha);
}
let shift = alpha - baseline_alpha.unwrap();
let verdict = struktura::HealthVerdict::from_shift(shift);
if json {
println!("{{\"alpha\":{:.4},\"r2\":{:.4},\"shift\":{:.4},\"verdict\":\"{}\"}}", alpha, r2, shift, verdict);
} else {
println!("alpha={:.4} R²={:.4} shift={:+.4} {}", alpha, r2, shift, verdict);
}
}
}
fn cmd_investigate(args: &[String]) {
let mut file_path = String::new();
let mut context_path = String::new();
let mut baseline_n = 0usize;
let mut json = false;
let mut i = 2;
while i < args.len() {
match args[i].as_str() {
"--context" if i + 1 < args.len() => { context_path = args[i + 1].clone(); i += 2; }
"--baseline" if i + 1 < args.len() => { baseline_n = args[i + 1].parse().unwrap_or(0); i += 2; }
"--json" => { json = true; i += 1; }
"--help" | "-h" => {
println!("struktura investigate <file.csv> [--context ctx.csv] [--baseline N] [--json]");
println!(" Run a full investigation: data quality, detection, incident grouping, report.");
println!(" --context Sidecar CSV with operating context (mode, commands, annotations)");
println!(" --baseline Calibration window in samples (default: first 1/3 of the file)");
println!(" --json Print incident records as JSON instead of human-readable report");
process::exit(0);
}
s if !s.starts_with('-') && file_path.is_empty() => { file_path = s.to_string(); i += 1; }
_ => { i += 1; }
}
}
if file_path.is_empty() {
eprintln!("Usage: struktura investigate <file.csv> [--context ctx.csv] [--baseline N]");
process::exit(2);
}
let content = std::fs::read_to_string(&file_path).unwrap_or_else(|e| {
eprintln!("cannot read {}: {}", file_path, e);
process::exit(2);
});
let (meas_names, meas_data, _schema, _header) = prepare_input(&content);
let channels = meas_data.len();
let samples = if channels > 0 { meas_data[0].len() } else { 0 };
if channels == 0 || samples < 64 {
eprintln!("need at least 1 channel and 64 samples, got {}ch x {}s", channels, samples);
process::exit(2);
}
let meas_names: Vec<&str> = meas_names.iter().map(|s| s.as_str()).collect();
let baseline = if baseline_n > 0 { baseline_n.min(meas_data[0].len()) } else { meas_data[0].len() / 3 };
let mut quality_issues = Vec::new();
for (ci, col) in meas_data.iter().enumerate() {
let bad = col.iter().filter(|v| !v.is_finite()).count();
if bad > 0 {
quality_issues.push(format!(" {}, {} non-finite values ({:.1}%)", meas_names[ci], bad, 100.0 * bad as f64 / col.len() as f64));
}
}
let mut timeline = struktura::context::ContextTimeline::new();
if !context_path.is_empty() {
match struktura::context::ContextTimeline::from_csv(&context_path, 0, 1, 2) {
Ok(t) => timeline = t,
Err(e) => eprintln!("warning: could not parse context CSV: {}", e),
}
}
eprintln!("struktura investigate: {} samples x {} channels, calibrated on {} rows", samples, meas_data.len(), baseline);
if !quality_issues.is_empty() {
eprintln!("\n⚠ Data quality issues:");
for q in &quality_issues { eprintln!("{}", q); }
eprintln!();
}
let (incidents, _monitor_export, _imputation) = match struktura::replay::run_investigation(&meas_data, baseline, &timeline) {
Ok(r) => r,
Err(e) => {
eprintln!("investigation failed: {}", e);
process::exit(2);
}
};
if json {
print!("[");
for (i, inc) in incidents.iter().enumerate() {
if i > 0 { print!(","); }
print!("{}", inc.to_json());
}
println!("]");
} else {
let report = struktura::report::investigation_report(&incidents, &meas_names, samples);
println!("{}", report);
}
if incidents.iter().any(|inc| !inc.evidence.is_empty()) {
process::exit(1);
}
}
fn cmd_case(args: &[String]) {
if args.len() < 3 {
eprintln!("Usage: struktura case <save|load> [options]");
process::exit(2);
}
match args[2].as_str() {
"save" => {
let mut recording_path = String::new();
let mut baseline_n = 0usize;
let mut name = String::new();
let mut out_dir = String::new();
let mut context_path = String::new();
let mut i = 3;
while i < args.len() {
match args[i].as_str() {
"--from" if i + 1 < args.len() => { recording_path = args[i + 1].clone(); i += 2; }
"--baseline" if i + 1 < args.len() => { baseline_n = args[i + 1].parse().unwrap_or(0); i += 2; }
"--name" if i + 1 < args.len() => { name = args[i + 1].clone(); i += 2; }
"--out" if i + 1 < args.len() => { out_dir = args[i + 1].clone(); i += 2; }
"--context" if i + 1 < args.len() => { context_path = args[i + 1].clone(); i += 2; }
"--help" | "-h" => {
println!("struktura case save --from <recording.csv> --name <case-name> [--baseline N] [--out <dir>] [--context ctx.csv]");
println!(" Save a recording and its investigation as a replayable case directory.");
println!(" --context Sidecar CSV with operating context (mode, commands, annotations)");
process::exit(0);
}
_ => { i += 1; }
}
}
if recording_path.is_empty() || name.is_empty() {
eprintln!("need --from <file> and --name <name>");
process::exit(2);
}
let content = std::fs::read_to_string(&recording_path).unwrap_or_else(|e| {
eprintln!("cannot read {}: {}", recording_path, e);
process::exit(2);
});
let input_hash = struktura::case::fingerprint_content(content.as_bytes());
let (meas_names, meas_data, column_schema, _header) = prepare_input(&content);
if meas_data.is_empty() || meas_data[0].is_empty() {
eprintln!("need at least 1 measurement channel and some samples");
process::exit(2);
}
let baseline = if baseline_n > 0 { baseline_n } else { meas_data[0].len() / 3 };
let mut timeline = struktura::context::ContextTimeline::new();
if !context_path.is_empty() {
match struktura::context::ContextTimeline::from_csv(&context_path, 0, 1, 2) {
Ok(t) => timeline = t,
Err(e) => eprintln!("warning: could not parse context CSV: {}", e),
}
}
let (incidents, monitor_export, imputation) = match struktura::replay::run_investigation(&meas_data, baseline, &timeline) {
Ok(r) => r,
Err(e) => { eprintln!("investigation failed: {}", e); process::exit(2); }
};
let dir = if out_dir.is_empty() {
std::path::PathBuf::from(format!("cases/{}", name))
} else {
std::path::PathBuf::from(&out_dir)
};
let nsamples = meas_data[0].len();
let rows: Vec<Vec<f64>> = (0..nsamples).map(|t| meas_data.iter().map(|c| c[t]).collect()).collect();
let config = struktura::case::CaseConfig { input_hash, monitor_export, column_schema, imputation };
match struktura::case::Case::save(&dir, &rows, &incidents, baseline, &name, &config, &timeline, &meas_names) {
Ok(c) => println!("case saved: {} ({} incidents, {} channels x {} samples)",
c.dir().display(), incidents.len(), meas_data.len(), nsamples),
Err(e) => { eprintln!("save failed: {}", e); process::exit(2); }
}
}
"load" => {
let dir = args.get(3).map(|s| s.as_str()).unwrap_or_else(|| {
eprintln!("Usage: struktura case load <dir>");
process::exit(2);
});
match struktura::case::Case::load(std::path::Path::new(dir)) {
Ok(c) => {
let m = c.manifest;
println!("case: {} (v{}, {} channels x {} samples, baseline {})",
m.name, m.detector_version, m.channels, m.samples, m.baseline_samples);
}
Err(e) => { eprintln!("load failed: {}", e); process::exit(2); }
}
}
other => {
eprintln!("unknown case subcommand: {} (try: save, load)", other);
process::exit(2);
}
}
}
fn cmd_replay(args: &[String]) {
let mut case_dir = String::new();
let mut baseline_override: Option<usize> = None;
let mut json = false;
let mut i = 2;
while i < args.len() {
match args[i].as_str() {
"--baseline" if i + 1 < args.len() => { baseline_override = args[i + 1].parse().ok(); i += 2; }
"--json" => { json = true; i += 1; }
"--help" | "-h" => {
println!("struktura replay <case-dir> [--baseline N] [--json]");
println!(" Re-run the detector on a saved case and diff against saved incidents.");
process::exit(0);
}
s if !s.starts_with('-') && case_dir.is_empty() => { case_dir = s.to_string(); i += 1; }
_ => { i += 1; }
}
}
if case_dir.is_empty() {
eprintln!("Usage: struktura replay <case-dir> [--baseline N] [--json]");
process::exit(2);
}
let case = match struktura::case::Case::load(std::path::Path::new(&case_dir)) {
Ok(c) => c,
Err(e) => { eprintln!("load failed: {}", e); process::exit(2); }
};
match struktura::replay::replay(&case, baseline_override) {
Ok((incidents, diff)) => {
if json {
println!("{}", diff.to_json());
} else {
let old_incidents = case.incidents().unwrap_or_default();
let rpt = struktura::report::replay_report(&diff, &old_incidents, &incidents);
println!("{}", rpt);
if !incidents.is_empty() {
println!();
let channel_names = case.channel_names();
let names: Vec<&str> = channel_names.iter().map(|s| s.as_str()).collect();
let rpt = struktura::report::investigation_report(&incidents, &names, case.manifest.samples);
println!("{}", rpt);
}
}
if incidents.iter().any(|inc| !inc.evidence.is_empty()) {
process::exit(1);
}
}
Err(e) => { eprintln!("replay failed: {}", e); process::exit(2); }
}
}
fn parse_csv_columns(content: &str) -> (Vec<String>, Vec<Vec<f64>>, Vec<usize>) {
let mut lines = content.lines();
let header_line = match lines.next() {
Some(h) => h,
None => return (vec![], vec![], vec![]),
};
let header: Vec<String> = header_line.split(',').map(|s| s.trim().to_string()).collect();
let ncols = header.len();
let mut cols: Vec<Vec<f64>> = (0..ncols).map(|_| Vec::new()).collect();
let mut invalid_rows = Vec::new();
for (row_idx, line) in lines.enumerate() {
let fields: Vec<&str> = line.split(',').collect();
if fields.len() == ncols {
for (i, f) in fields.iter().enumerate() {
cols[i].push(f.trim().parse::<f64>().unwrap_or(f64::NAN));
}
} else {
invalid_rows.push(row_idx);
for col in cols.iter_mut() {
col.push(f64::NAN);
}
}
}
(header, cols, invalid_rows)
}
fn prepare_input(content: &str) -> (Vec<String>, Vec<Vec<f64>>, struktura::context::ColumnSchema, Vec<String>) {
let (header, cols, invalid_rows) = parse_csv_columns(content);
if !invalid_rows.is_empty() {
eprintln!(
"warning: {} malformed row(s) in input (rows {:?}, set to NaN across all channels)",
invalid_rows.len(),
invalid_rows
);
}
let schema = struktura::context::ColumnSchema::from_header(
&header.iter().map(|s| s.as_str()).collect::<Vec<_>>()
);
let measurement_cols = schema.measurement_indices();
let meas_cols: Vec<Vec<f64>> = measurement_cols.iter().map(|&i| cols[i].clone()).collect();
let meas_names: Vec<String> = measurement_cols.iter().map(|&i| header[i].clone()).collect();
(meas_names, meas_cols, schema, header)
}
fn cmd_copilot_compare(args: &[String]) {
use struktura::monitor::HybridMonitor;
use struktura::autopilot::{AutoPilot, Event};
let file_path = args.get(2).unwrap_or_else(|| {
eprintln!("Usage: struktura copilot-compare <file.csv> [--col N]");
eprintln!(" Side-by-side: DFA structural monitor vs boolean amplitude threshold.");
eprintln!(" Prints the first alarm row of each; run a healthy-period control before calling either a detection.");
process::exit(1);
});
let content = std::fs::read_to_string(file_path).unwrap_or_else(|e| {
eprintln!("cannot read {}: {}", file_path, e);
process::exit(2);
});
let parsed = parse_multi_csv(&content);
let rows = parsed.rows;
if rows.is_empty() { eprintln!("no numeric rows"); process::exit(2); }
let n = rows.len();
let ncols = rows[0].len();
let calib_n = (n / 3).clamp(192, 100_000).min(n);
let channels: Vec<Vec<f64>> = (0..ncols)
.map(|ch| rows.iter().map(|r| r.get(ch).copied().unwrap_or(0.0)).collect())
.collect();
let mut thresholds = vec![0.0f64; ncols];
for ch in 0..ncols {
let mut calib_vals: Vec<f64> = channels[ch][..calib_n].iter().map(|v| v.abs()).collect();
calib_vals.sort_by(|a, b| a.partial_cmp(b).unwrap_or(std::cmp::Ordering::Equal));
let p95_idx = (calib_vals.len() as f64 * 0.95) as usize;
thresholds[ch] = calib_vals.get(p95_idx.min(calib_vals.len() - 1)).copied().unwrap_or(1.0) * 1.5;
}
let calib: Vec<Vec<f64>> = channels.iter().map(|c| c[..calib_n].to_vec()).collect();
let mon = match HybridMonitor::calibrate(&calib) {
Some(m) => m,
None => { eprintln!("calibration failed (need >= 192 samples)"); process::exit(2); }
};
let mut ap = AutoPilot::new(mon);
let valid: Vec<bool> = vec![true; ncols];
let mut sample = vec![0.0f64; ncols];
let mut first_dfa_alarm: Option<usize> = None;
let mut first_confirmed: Option<usize> = None;
let mut first_bool_alarm: Option<usize> = None;
let mut last_alarm_leg: Vec<(usize, u8)> = Vec::new();
println!();
println!(" \x1b[1mCopilot Boolean vs Struktura DFA: {}\x1b[0m", file_path);
println!(" {} samples × {} channels, calibrated on {} rows", n, ncols, calib_n);
println!(" ─────────────────────────────────────────────────────────────────────");
println!(" {:>6} {:>10} {:>14} {:>20}", "row", "amplitude", "bool monitor", "struktura DFA");
println!(" {:>6} {:>10} {:>14} {:>20}", "───", "─────────", "────────────", "─────────────");
let baseline_law = analyze(&channels[0][..calib_n]);
println!(" {:>6} {:>10} {:>14} α={:<6.3} baseline", calib_n, "in-spec", "✓ OK", baseline_law.dfa.alpha);
for t in calib_n..n {
for ch in 0..ncols { sample[ch] = channels[ch][t]; }
let max_amp = (0..ncols).map(|ch| channels[ch][t].abs()).fold(0.0f64, f64::max);
let bool_tripped = (0..ncols).any(|ch| channels[ch][t].abs() > thresholds[ch]);
if bool_tripped && first_bool_alarm.is_none() { first_bool_alarm = Some(t); }
let mut dfa_event = None;
for ev in ap.push(&sample, &valid) {
match &ev {
Event::Alarm { report, .. } => {
let leg_id = report.leg as u8;
let dup = last_alarm_leg.iter().any(|&(lt, ll)| ll == leg_id && t.saturating_sub(lt) < 50);
last_alarm_leg.retain(|&(lt, _)| t.saturating_sub(lt) < 50);
last_alarm_leg.push((t, leg_id));
if !dup {
let explanation = struktura::monitor::explain_alarm(report);
dfa_event = Some(format!("⚠ {}", explanation.chars().take(30).collect::<String>()));
if first_dfa_alarm.is_none() { first_dfa_alarm = Some(t); }
}
}
Event::AdaptationStarted { .. } => {
dfa_event = Some("↻ learning new baseline".into());
}
Event::RolledBack { .. } => {
dfa_event = Some("✗ FAULT CONFIRMED".into());
if first_confirmed.is_none() { first_confirmed = Some(t); }
}
Event::Recalibrated { .. } => {
dfa_event = Some("✓ new baseline accepted".into());
}
_ => {}
}
}
let bool_just_tripped = bool_tripped && first_bool_alarm == Some(t);
let is_dfa_event = dfa_event.is_some();
if is_dfa_event || bool_just_tripped {
let bool_str = if bool_tripped { "\x1b[31m✗ ALARM\x1b[0m" } else { "✓ OK" };
let dfa_str = dfa_event.unwrap_or_else(|| "-".into());
let amp_str = if bool_tripped { format!("\x1b[31m{:.4}\x1b[0m", max_amp) } else { format!("{:.4}", max_amp) };
println!(" {:>6} {:>10} {:>14} {}", t, amp_str, bool_str, dfa_str);
}
}
println!(" ─────────────────────────────────────────────────────────────────────");
println!();
let row = |r: Option<usize>| r.map_or_else(|| "none".to_string(), |t| t.to_string());
println!(" first struktura monitor alarm (any leg): row {}", row(first_dfa_alarm));
println!(" first confirmed fault (guarded rollback): row {}", row(first_confirmed));
println!(" first amplitude-threshold trip: row {}", row(first_bool_alarm));
println!(" threshold = 1.5 x the 95th percentile of |x| in the calibration rows, per channel");
println!();
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn parse_csv_columns_marks_malformed_rows_entirely_invalid_preserving_row_count() {
let csv = "time,a,b\n1,10,20\n2,11\n3,12,13,extra\n4,14,15\n";
let (header, cols, invalid_rows) = parse_csv_columns(csv);
assert_eq!(header, vec!["time", "a", "b"]);
assert_eq!(cols[0].len(), 4);
assert_eq!(cols[1].len(), 4);
assert_eq!(cols[2].len(), 4);
assert_eq!(invalid_rows, vec![1, 2]);
assert!(cols[0][1].is_nan() && cols[1][1].is_nan() && cols[2][1].is_nan());
assert!(cols[0][2].is_nan() && cols[1][2].is_nan() && cols[2][2].is_nan());
assert_eq!(cols[0][0], 1.0);
assert_eq!(cols[1][0], 10.0);
assert_eq!(cols[0][3], 4.0);
assert_eq!(cols[2][3], 15.0);
}
#[test]
fn parse_csv_columns_treats_unparseable_field_as_nan_without_flagging_whole_row_invalid() {
let csv = "a,mode\n1,cruise\n2,idle\n";
let (_, cols, invalid_rows) = parse_csv_columns(csv);
assert!(invalid_rows.is_empty());
assert_eq!(cols[0], vec![1.0, 2.0]);
assert!(cols[1][0].is_nan() && cols[1][1].is_nan());
}
#[test]
fn parse_csv_columns_reports_zero_malformed_on_clean_input() {
let (_, _, invalid_rows) = parse_csv_columns("a,b\n1,2\n3,4\n");
assert!(invalid_rows.is_empty());
}
#[test]
fn prepare_input_filters_to_measurement_columns_only() {
let csv = "time,mode,motor_current\n0,cruise,1.0\n1,cruise,1.1\n2,idle,1.2\n";
let (names, cols, schema, header) = prepare_input(csv);
assert_eq!(names, vec!["motor_current"]);
assert_eq!(cols.len(), 1);
assert_eq!(cols[0], vec![1.0, 1.1, 1.2]);
assert_eq!(header, vec!["time", "mode", "motor_current"]);
assert_eq!(schema.measurement_indices(), vec![2]);
}
}