///|
/// A lightweight scalar telemetry record for dashboards and field logs.
pub struct TelemetrySample {
timestamp : Int
channel : String
value : Double
quality : Double
result : UpdateResult
} derive(Debug)
///|
pub fn TelemetrySample::new(
timestamp : Int,
channel : String,
value : Double,
quality : Double,
result : UpdateResult,
) -> TelemetrySample {
{
timestamp,
channel,
value,
quality: if quality < 0.0 {
0.0
} else if quality > 1.0 {
1.0
} else {
quality
},
result,
}
}
///|
pub fn TelemetrySample::timestamp(self : TelemetrySample) -> Int {
self.timestamp
}
///|
pub fn TelemetrySample::channel(self : TelemetrySample) -> String {
self.channel
}
///|
pub fn TelemetrySample::value(self : TelemetrySample) -> Double {
self.value
}
///|
pub fn TelemetrySample::quality(self : TelemetrySample) -> Double {
self.quality
}
///|
pub fn TelemetrySample::result(self : TelemetrySample) -> UpdateResult {
self.result
}
///|
pub fn TelemetrySample::is_finite(self : TelemetrySample) -> Bool {
!self.value.is_nan() && !self.value.is_inf()
}
///|
pub struct TelemetrySeries {
samples : Array[TelemetrySample]
capacity : Int
mut rejected : Int
}
///|
pub fn TelemetrySeries::new(capacity : Int) -> TelemetrySeries {
{
samples: [],
capacity: if capacity < 0 {
0
} else {
capacity
},
rejected: 0,
}
}
///|
pub fn TelemetrySeries::push(
self : TelemetrySeries,
sample : TelemetrySample,
) -> Bool {
if self.capacity == 0 || !sample.is_finite() {
self.rejected = self.rejected + 1
return false
}
if self.samples.length() > 0 &&
sample.timestamp() < self.samples[self.samples.length() - 1].timestamp() {
self.rejected = self.rejected + 1
return false
}
if self.samples.length() >= self.capacity {
for i in 1.. Array[TelemetrySample] {
self.samples.copy()
}
///|
pub fn TelemetrySeries::length(self : TelemetrySeries) -> Int {
self.samples.length()
}
///|
pub fn TelemetrySeries::rejected(self : TelemetrySeries) -> Int {
self.rejected
}
///|
pub fn TelemetrySeries::mean(self : TelemetrySeries) -> Double {
let stats = RunningStats::new()
for sample in self.samples {
stats.add(sample.value())
}
stats.mean()
}
///|
pub fn TelemetrySeries::variance(self : TelemetrySeries) -> Double {
let stats = RunningStats::new()
for sample in self.samples {
stats.add(sample.value())
}
stats.variance()
}
///|
pub fn TelemetrySeries::quality(self : TelemetrySeries) -> Double {
if self.samples.length() == 0 {
return 0.0
}
let mut total = 0.0
for sample in self.samples {
total = total + sample.quality()
}
total / self.samples.length().to_double()
}
///|
pub fn TelemetrySeries::duration(self : TelemetrySeries) -> Int {
if self.samples.length() < 2 {
0
} else {
self.samples[self.samples.length() - 1].timestamp() -
self.samples[0].timestamp()
}
}
///|
pub fn TelemetrySeries::reset(self : TelemetrySeries) -> Unit {
self.samples.clear()
self.rejected = 0
}
///|
pub fn telemetry_from_filter(
timestamp : Int,
channel : String,
filter : Kalman1D,
result : UpdateResult,
) -> TelemetrySample {
let quality = if result is Accepted {
1.0
} else if result is MissingMeasurement {
0.5
} else {
0.0
}
TelemetrySample::new(timestamp, channel, filter.state(), quality, result)
}
///|
pub fn telemetry_to_csv(samples : Array[TelemetrySample]) -> String {
let mut result = "timestamp,channel,value,quality,result\n"
for sample in samples {
result = result +
sample.timestamp().to_string() +
"," +
sample.channel() +
"," +
sample.value().to_string() +
"," +
sample.quality().to_string() +
"," +
update_result_to_string(sample.result()) +
"\n"
}
result
}
///|
pub fn telemetry_quality_histogram(
samples : Array[TelemetrySample],
buckets : Int,
) -> Histogram {
let histogram = Histogram::new(0.0, 1.0, buckets)
for sample in samples {
histogram.add(sample.quality()) |> ignore
}
histogram
}
///|
pub fn telemetry_outlier_count(
samples : Array[TelemetrySample],
threshold : Double,
) -> Int {
let limit = if threshold < 0.0 { 0.0 } else { threshold }
let mut count = 0
for sample in samples {
if sample.value().abs() > limit {
count = count + 1
}
}
count
}
///|
pub fn telemetry_merge(
left : Array[TelemetrySample],
right : Array[TelemetrySample],
) -> Array[TelemetrySample] {
let result : Array[TelemetrySample] = []
for sample in left {
result.push(sample)
}
for sample in right {
result.push(sample)
}
for i in 1.. 0 && result[index - 1].timestamp() > value.timestamp() {
result[index] = result[index - 1]
index = index - 1
}
result[index] = value
}
result
}
///|
pub fn telemetry_resample(
samples : Array[TelemetrySample],
timestamps : Array[Int],
) -> Array[TelemetrySample] {
let result : Array[TelemetrySample] = []
if samples.length() == 0 {
return result
}
for timestamp in timestamps {
let mut best = samples[0]
let mut distance = (timestamp - best.timestamp()).abs()
for sample in samples {
let current = (timestamp - sample.timestamp()).abs()
if current < distance {
best = sample
distance = current
}
}
result.push(
TelemetrySample::new(
timestamp,
best.channel(),
best.value(),
best.quality(),
best.result(),
),
)
}
result
}
///|
pub fn telemetry_status(samples : Array[TelemetrySample]) -> FilterStatus {
if samples.length() == 0 {
WarmingUp
} else {
let mut failures = 0
for sample in samples {
match sample.result() {
Accepted => ()
RejectedByGate
| InvalidMeasurement
| SingularInnovation
| MissingMeasurement => failures = failures + 1
}
}
if failures == 0 {
Healthy
} else if failures >= samples.length() / 2 {
Faulted
} else {
Degraded
}
}
}
///|
pub struct TelemetrySummary {
samples : Int
accepted : Int
missing : Int
rejected : Int
mean : Double
variance : Double
quality : Double
status : FilterStatus
} derive(Debug)
///|
pub fn TelemetrySummary::samples(self : TelemetrySummary) -> Int {
self.samples
}
///|
pub fn TelemetrySummary::accepted(self : TelemetrySummary) -> Int {
self.accepted
}
///|
pub fn TelemetrySummary::missing(self : TelemetrySummary) -> Int {
self.missing
}
///|
pub fn TelemetrySummary::rejected(self : TelemetrySummary) -> Int {
self.rejected
}
///|
pub fn TelemetrySummary::mean(self : TelemetrySummary) -> Double {
self.mean
}
///|
pub fn TelemetrySummary::variance(self : TelemetrySummary) -> Double {
self.variance
}
///|
pub fn TelemetrySummary::quality(self : TelemetrySummary) -> Double {
self.quality
}
///|
pub fn TelemetrySummary::status(self : TelemetrySummary) -> FilterStatus {
self.status
}
///|
pub fn summarize_telemetry(
samples : Array[TelemetrySample],
) -> TelemetrySummary {
let stats = RunningStats::new()
let mut accepted = 0
let mut missing = 0
let mut rejected = 0
let mut quality = 0.0
for sample in samples {
stats.add(sample.value())
quality = quality + sample.quality()
match sample.result() {
Accepted => accepted = accepted + 1
MissingMeasurement => missing = missing + 1
RejectedByGate | InvalidMeasurement | SingularInnovation =>
rejected = rejected + 1
}
}
let count = samples.length()
{
samples: count,
accepted,
missing,
rejected,
mean: stats.mean(),
variance: stats.variance(),
quality: if count == 0 {
0.0
} else {
quality / count.to_double()
},
status: telemetry_status(samples),
}
}
///|
pub fn telemetry_gap_count(
samples : Array[TelemetrySample],
expected_period : Int,
) -> Int {
if samples.length() < 2 || expected_period < 1 {
return 0
}
let mut gaps = 0
for i in 1.. expected_period {
gaps = gaps + 1
}
}
gaps
}
///|
pub fn telemetry_range(samples : Array[TelemetrySample]) -> Double {
if samples.length() == 0 {
return 0.0
}
let mut minimum = samples[0].value()
let mut maximum = minimum
for sample in samples {
if sample.value() < minimum {
minimum = sample.value()
}
if sample.value() > maximum {
maximum = sample.value()
}
}
maximum - minimum
}
///|
pub fn telemetry_all_finite(samples : Array[TelemetrySample]) -> Bool {
for sample in samples {
if !sample.is_finite() {
return false
}
}
true
}
///|
pub fn telemetry_status_score(status : FilterStatus) -> Double {
match status {
Healthy => 1.0
WarmingUp => 0.75
Degraded => 0.5
Faulted => 0.0
}
}
///|
pub fn telemetry_acceptance_rate(summary : TelemetrySummary) -> Double {
if summary.samples() == 0 {
0.0
} else {
summary.accepted().to_double() / summary.samples().to_double()
}
}
///|
pub fn telemetry_failure_rate(summary : TelemetrySummary) -> Double {
if summary.samples() == 0 {
0.0
} else {
(summary.missing() + summary.rejected()).to_double() /
summary.samples().to_double()
}
}
///|
pub fn telemetry_is_usable(
summary : TelemetrySummary,
minimum_quality : Double,
) -> Bool {
let healthy = match summary.status() {
Faulted => false
WarmingUp | Healthy | Degraded => true
}
summary.samples() > 0 && summary.quality() >= minimum_quality && healthy
}
///|
pub fn telemetry_quality_adjusted_value(
sample : TelemetrySample,
fallback : Double,
) -> Double {
if !sample.is_finite() {
fallback
} else {
sample.value() * sample.quality() + fallback * (1.0 - sample.quality())
}
}
///|
pub fn telemetry_values(samples : Array[TelemetrySample]) -> Array[Double] {
samples.map(sample => sample.value())
}
///|
pub fn telemetry_quality_weighted_mean(
samples : Array[TelemetrySample],
) -> Double {
let mut total = 0.0
let mut weight = 0.0
for sample in samples {
total = total + sample.value() * sample.quality()
weight = weight + sample.quality()
}
if weight <= 0.0 {
0.0
} else {
total / weight
}
}
///|
pub fn telemetry_sample_count(
samples : Array[TelemetrySample],
result : UpdateResult,
) -> Int {
let mut count = 0
for sample in samples {
if sample.result() == result {
count = count + 1
}
}
count
}
///|
pub fn telemetry_is_monotonic(samples : Array[TelemetrySample]) -> Bool {
for i in 1..