///|
/// A timestamped vector sample used by asynchronous sensor adapters.
pub struct TimedVector {
timestamp : Int
values : Array[Double]
quality : Double
} derive(Debug)
///|
/// Construct a timestamped vector. The values are copied at the boundary.
pub fn TimedVector::new(
timestamp : Int,
values : Array[Double],
quality : Double,
) -> TimedVector {
{ timestamp, values: values.copy(), quality: quality.clamp(min=0.0, max=1.0) }
}
///|
/// Return the timestamp.
pub fn TimedVector::timestamp(self : TimedVector) -> Int {
self.timestamp
}
///|
/// Return a copy of the sample vector.
pub fn TimedVector::values(self : TimedVector) -> Array[Double] {
self.values.copy()
}
///|
/// Return the sample quality.
pub fn TimedVector::quality(self : TimedVector) -> Double {
self.quality
}
///|
/// Return the vector dimension.
pub fn TimedVector::dimension(self : TimedVector) -> Int {
self.values.length()
}
///|
/// Return whether the sample is finite and non-empty.
pub fn TimedVector::is_valid(self : TimedVector) -> Bool {
self.values.length() > 0 &&
vector_is_finite(self.values) &&
!self.quality.is_nan() &&
!self.quality.is_inf()
}
///|
/// Replace the quality without changing the vector.
pub fn TimedVector::with_quality(
self : TimedVector,
quality : Double,
) -> TimedVector {
TimedVector::new(self.timestamp, self.values, quality)
}
///|
/// Shift the timestamp by an integer offset.
pub fn TimedVector::shift(self : TimedVector, offset : Int) -> TimedVector {
TimedVector::new(self.timestamp + offset, self.values, self.quality)
}
///|
/// Scale all vector components.
pub fn TimedVector::scale(self : TimedVector, factor : Double) -> TimedVector {
TimedVector::new(
self.timestamp,
vector_scale(self.values, factor),
self.quality,
)
}
///|
/// A policy controlling interpolation and gap handling.
pub struct AlignmentPolicy {
tolerance : Int
max_gap : Int
allow_extrapolation : Bool
minimum_quality : Double
} derive(Debug)
///|
/// Construct an alignment policy.
pub fn AlignmentPolicy::new(
tolerance : Int,
max_gap : Int,
allow_extrapolation : Bool,
minimum_quality : Double,
) -> AlignmentPolicy {
{
tolerance: if tolerance < 0 {
0
} else {
tolerance
},
max_gap: if max_gap < 0 {
0
} else {
max_gap
},
allow_extrapolation,
minimum_quality: minimum_quality.clamp(min=0.0, max=1.0),
}
}
///|
/// Return a permissive policy for embedded sensor streams.
pub fn permissive_alignment_policy() -> AlignmentPolicy {
AlignmentPolicy::new(5, 1000, false, 0.0)
}
///|
/// Return a strict policy for synchronized control loops.
pub fn strict_alignment_policy() -> AlignmentPolicy {
AlignmentPolicy::new(0, 100, false, 0.5)
}
///|
/// Return timestamp tolerance.
pub fn AlignmentPolicy::tolerance(self : AlignmentPolicy) -> Int {
self.tolerance
}
///|
/// Return maximum interpolation gap.
pub fn AlignmentPolicy::max_gap(self : AlignmentPolicy) -> Int {
self.max_gap
}
///|
/// Return extrapolation policy.
pub fn AlignmentPolicy::allow_extrapolation(self : AlignmentPolicy) -> Bool {
self.allow_extrapolation
}
///|
/// Return minimum acceptable quality.
pub fn AlignmentPolicy::minimum_quality(self : AlignmentPolicy) -> Double {
self.minimum_quality
}
///|
/// The outcome class of a timestamp alignment query.
pub(all) enum AlignmentStatus {
Exact
Nearest
Interpolated
Extrapolated
Missing
Invalid
} derive(Debug, Eq)
///|
/// A value aligned to a requested timestamp.
pub struct AlignedVector {
timestamp : Int
values : Array[Double]
status : AlignmentStatus
source_timestamp : Int
quality : Double
distance : Int
} derive(Debug)
///|
/// Construct an aligned vector.
pub fn AlignedVector::new(
timestamp : Int,
values : Array[Double],
status : AlignmentStatus,
source_timestamp : Int,
quality : Double,
distance : Int,
) -> AlignedVector {
{
timestamp,
values: values.copy(),
status,
source_timestamp,
quality,
distance,
}
}
///|
/// Return target timestamp.
pub fn AlignedVector::timestamp(self : AlignedVector) -> Int {
self.timestamp
}
///|
/// Return aligned values.
pub fn AlignedVector::values(self : AlignedVector) -> Array[Double] {
self.values.copy()
}
///|
/// Return status.
pub fn AlignedVector::status(self : AlignedVector) -> AlignmentStatus {
self.status
}
///|
/// Return source timestamp.
pub fn AlignedVector::source_timestamp(self : AlignedVector) -> Int {
self.source_timestamp
}
///|
/// Return propagated quality.
pub fn AlignedVector::quality(self : AlignedVector) -> Double {
self.quality
}
///|
/// Return absolute source/target distance.
pub fn AlignedVector::distance(self : AlignedVector) -> Int {
self.distance
}
///|
/// Return whether alignment produced a value.
pub fn AlignedVector::is_present(self : AlignedVector) -> Bool {
match self.status {
Missing | Invalid => false
_ => true
}
}
///|
/// A summary of an alignment pass.
pub struct AlignmentReport {
requested : Int
produced : Int
exact : Int
nearest : Int
interpolated : Int
extrapolated : Int
missing : Int
invalid : Int
maximum_distance : Int
mean_distance : Double
} derive(Debug)
///|
/// Return requested count.
pub fn AlignmentReport::requested(self : AlignmentReport) -> Int {
self.requested
}
///|
/// Return produced count.
pub fn AlignmentReport::produced(self : AlignmentReport) -> Int {
self.produced
}
///|
/// Return exact count.
pub fn AlignmentReport::exact(self : AlignmentReport) -> Int {
self.exact
}
///|
/// Return nearest count.
pub fn AlignmentReport::nearest(self : AlignmentReport) -> Int {
self.nearest
}
///|
/// Return interpolated count.
pub fn AlignmentReport::interpolated(self : AlignmentReport) -> Int {
self.interpolated
}
///|
/// Return extrapolated count.
pub fn AlignmentReport::extrapolated(self : AlignmentReport) -> Int {
self.extrapolated
}
///|
/// Return missing count.
pub fn AlignmentReport::missing(self : AlignmentReport) -> Int {
self.missing
}
///|
/// Return invalid count.
pub fn AlignmentReport::invalid(self : AlignmentReport) -> Int {
self.invalid
}
///|
/// Return maximum timestamp distance.
pub fn AlignmentReport::maximum_distance(self : AlignmentReport) -> Int {
self.maximum_distance
}
///|
/// Return mean timestamp distance.
pub fn AlignmentReport::mean_distance(self : AlignmentReport) -> Double {
self.mean_distance
}
///|
/// Return the produced fraction.
pub fn AlignmentReport::coverage(self : AlignmentReport) -> Double {
if self.requested == 0 {
0.0
} else {
self.produced.to_double() / self.requested.to_double()
}
}
///|
/// Return the fraction produced by interpolation or exact matching.
pub fn AlignmentReport::usable_coverage(self : AlignmentReport) -> Double {
if self.requested == 0 {
0.0
} else {
(self.exact + self.interpolated).to_double() / self.requested.to_double()
}
}
///|
/// A linear clock correction estimated from paired timestamps.
pub struct SensorClockFit {
offset : Double
drift : Double
reference : Int
residual_rms : Double
samples : Int
valid : Bool
} derive(Debug)
///|
/// Return clock offset.
pub fn SensorClockFit::offset(self : SensorClockFit) -> Double {
self.offset
}
///|
/// Return clock drift as a fractional rate.
pub fn SensorClockFit::drift(self : SensorClockFit) -> Double {
self.drift
}
///|
/// Return reference timestamp.
pub fn SensorClockFit::reference(self : SensorClockFit) -> Int {
self.reference
}
///|
/// Return residual RMS.
pub fn SensorClockFit::residual_rms(self : SensorClockFit) -> Double {
self.residual_rms
}
///|
/// Return pair count.
pub fn SensorClockFit::samples(self : SensorClockFit) -> Int {
self.samples
}
///|
/// Return whether the fit is usable.
pub fn SensorClockFit::valid(self : SensorClockFit) -> Bool {
self.valid
}
///|
/// Correct a timestamp using this clock fit.
pub fn SensorClockFit::correct(self : SensorClockFit, timestamp : Int) -> Int {
let delta = (timestamp - self.reference).to_double()
(timestamp.to_double() + self.offset + self.drift * delta).to_int()
}
///|
/// Return a copy of a series sorted by timestamp using insertion sort.
pub fn sort_timed_vectors(samples : Array[TimedVector]) -> Array[TimedVector] {
let result = samples.copy()
for i in 1.. 0 && result[j - 1].timestamp() > current.timestamp() {
result[j] = result[j - 1]
j = j - 1
}
result[j] = current
}
result
}
///|
/// Return whether timestamps are monotonic.
pub fn timed_vectors_are_monotonic(samples : Array[TimedVector]) -> Bool {
for i in 1.. Array[TimedVector] {
let result = []
for sample in samples {
if sample.is_valid() && sample.quality() >= minimum_quality {
result.push(sample)
}
}
result
}
///|
/// Find the nearest sample to a timestamp.
pub fn timed_vector_nearest(
samples : Array[TimedVector],
timestamp : Int,
policy : AlignmentPolicy,
) -> AlignedVector {
if samples.length() == 0 {
return AlignedVector::new(timestamp, [], Missing, timestamp, 0.0, 0)
}
let mut best = -1
let mut distance = 0
for i, sample in samples {
let current = (sample.timestamp() - timestamp).abs()
if best < 0 || current < distance {
best = i
distance = current
}
}
let sample = samples[best]
if !sample.is_valid() ||
sample.quality() < policy.minimum_quality() ||
distance > policy.tolerance() {
AlignedVector::new(
timestamp,
[],
Missing,
sample.timestamp(),
sample.quality(),
distance,
)
} else if distance == 0 {
AlignedVector::new(
timestamp,
sample.values(),
Exact,
sample.timestamp(),
sample.quality(),
distance,
)
} else {
AlignedVector::new(
timestamp,
sample.values(),
Nearest,
sample.timestamp(),
sample.quality(),
distance,
)
}
}
///|
/// Interpolate a vector between two samples.
pub fn interpolate_timed_vectors(
left : TimedVector,
right : TimedVector,
timestamp : Int,
) -> AlignedVector {
let span = right.timestamp() - left.timestamp()
let distance = (timestamp - left.timestamp()).abs()
if !left.is_valid() ||
!right.is_valid() ||
left.dimension() != right.dimension() {
return AlignedVector::new(
timestamp,
[],
Invalid,
left.timestamp(),
0.0,
distance,
)
}
if span == 0 {
return AlignedVector::new(
timestamp,
left.values(),
Exact,
left.timestamp(),
left.quality(),
distance,
)
}
let amount = (timestamp - left.timestamp()).to_double() / span.to_double()
AlignedVector::new(
timestamp,
vector_lerp(left.values(), right.values(), amount),
Interpolated,
left.timestamp(),
left.quality().min(right.quality()),
if amount < 0.0 {
(left.timestamp() - timestamp).abs()
} else if amount > 1.0 {
(timestamp - right.timestamp()).abs()
} else {
0
},
)
}
///|
/// Interpolate a series without extrapolating outside its support.
pub fn timed_vector_interpolate(
samples : Array[TimedVector],
timestamp : Int,
policy : AlignmentPolicy,
) -> AlignedVector {
if samples.length() == 0 {
return AlignedVector::new(timestamp, [], Missing, timestamp, 0.0, 0)
}
for sample in samples {
if sample.timestamp() == timestamp {
return timed_vector_nearest(samples, timestamp, policy)
}
}
if timestamp < samples[0].timestamp() ||
timestamp > samples[samples.length() - 1].timestamp() {
if !policy.allow_extrapolation() {
return AlignedVector::new(timestamp, [], Missing, timestamp, 0.0, 0)
}
if timestamp < samples[0].timestamp() && samples.length() >= 2 {
let result = interpolate_timed_vectors(samples[0], samples[1], timestamp)
return AlignedVector::new(
timestamp,
result.values(),
Extrapolated,
result.source_timestamp(),
result.quality(),
result.distance(),
)
}
if samples.length() >= 2 {
let last = samples.length() - 1
let result = interpolate_timed_vectors(
samples[last - 1],
samples[last],
timestamp,
)
return AlignedVector::new(
timestamp,
result.values(),
Extrapolated,
result.source_timestamp(),
result.quality(),
result.distance(),
)
}
return AlignedVector::new(timestamp, [], Missing, timestamp, 0.0, 0)
}
for i in 1.. policy.max_gap() {
return AlignedVector::new(
timestamp,
[],
Missing,
timestamp,
0.0,
result.distance(),
)
}
return result
}
}
AlignedVector::new(timestamp, [], Missing, timestamp, 0.0, 0)
}
///|
/// Align a series to an explicit timestamp grid.
pub fn align_timed_vectors(
samples : Array[TimedVector],
timestamps : Array[Int],
policy : AlignmentPolicy,
) -> (Array[AlignedVector], AlignmentReport) {
let aligned = []
let mut exact = 0
let mut nearest = 0
let mut interpolated = 0
let mut extrapolated = 0
let mut missing = 0
let mut invalid = 0
let mut maximum_distance = 0
let mut distance_sum = 0
for timestamp in timestamps {
let result = timed_vector_interpolate(samples, timestamp, policy)
aligned.push(result)
if result.distance() > maximum_distance {
maximum_distance = result.distance()
}
distance_sum = distance_sum + result.distance()
match result.status() {
Exact => exact = exact + 1
Nearest => nearest = nearest + 1
Interpolated => interpolated = interpolated + 1
Extrapolated => extrapolated = extrapolated + 1
Missing => missing = missing + 1
Invalid => invalid = invalid + 1
}
}
let produced = exact + nearest + interpolated + extrapolated
let mean_distance = if timestamps.length() == 0 {
0.0
} else {
distance_sum.to_double() / timestamps.length().to_double()
}
(
aligned,
{
requested: timestamps.length(),
produced,
exact,
nearest,
interpolated,
extrapolated,
missing,
invalid,
maximum_distance,
mean_distance,
},
)
}
///|
/// Build a regular timestamp grid.
pub fn regular_timestamps(start : Int, end : Int, period : Int) -> Array[Int] {
if period <= 0 || end < start {
return []
}
let result = []
let mut timestamp = start
while timestamp <= end {
result.push(timestamp)
timestamp = timestamp + period
}
result
}
///|
/// Return timestamp gaps larger than an expected period.
pub fn timed_vector_gaps(
samples : Array[TimedVector],
expected_period : Int,
) -> Array[Int] {
let result = []
if expected_period <= 0 {
return result
}
for i in 1.. expected_period {
result.push(gap)
}
}
result
}
///|
/// Return the median timestamp interval.
pub fn timed_vector_median_period(samples : Array[TimedVector]) -> Int {
if samples.length() < 2 {
return 0
}
let periods = []
for i in 1.. SensorClockFit {
let count = if local_timestamps.length() < reference_timestamps.length() {
local_timestamps.length()
} else {
reference_timestamps.length()
}
if count == 0 {
return {
offset: 0.0,
drift: 0.0,
reference: 0,
residual_rms: 0.0,
samples: 0,
valid: false,
}
}
let reference = local_timestamps[0]
let mut sx = 0.0
let mut sy = 0.0
let mut sxx = 0.0
let mut sxy = 0.0
for i in 0..= 2,
}
}
///|
/// Apply a clock fit to a timestamped series.
pub fn correct_sensor_clock(
samples : Array[TimedVector],
fit : SensorClockFit,
) -> Array[TimedVector] {
Array::makei(samples.length(), i => {
TimedVector::new(
fit.correct(samples[i].timestamp()),
samples[i].values(),
samples[i].quality(),
)
})
}
///|
/// Correct a series using the clock fit while retaining vector and quality.
pub fn correct_sensor_clock_values(
samples : Array[TimedVector],
fit : SensorClockFit,
) -> Array[TimedVector] {
let result = []
for sample in samples {
let corrected = fit.correct(sample.timestamp())
result.push(TimedVector::new(corrected, sample.values(), sample.quality()))
}
result
}
///|
/// Merge two timestamped streams and retain deterministic timestamp ordering.
pub fn merge_timed_vectors(
left : Array[TimedVector],
right : Array[TimedVector],
) -> Array[TimedVector] {
let result = left.copy()
for sample in right {
result.push(sample)
}
sort_timed_vectors(result)
}
///|
/// Merge multiple streams.
pub fn merge_timed_vector_streams(
streams : Array[Array[TimedVector]],
) -> Array[TimedVector] {
let result = []
for stream in streams {
for sample in stream {
result.push(sample)
}
}
sort_timed_vectors(result)
}
///|
/// Remove duplicate timestamps, retaining the higher quality sample.
pub fn deduplicate_timed_vectors(
samples : Array[TimedVector],
) -> Array[TimedVector] {
let sorted = sort_timed_vectors(samples)
let result : Array[TimedVector] = []
for sample in sorted {
if result.length() == 0 ||
sample.timestamp() != result[result.length() - 1].timestamp() {
result.push(sample)
} else if sample.quality() > result[result.length() - 1].quality() {
result[result.length() - 1] = sample
}
}
result
}
///|
/// Resample a stream on a regular grid and return both values and report.
pub fn resample_timed_vectors(
samples : Array[TimedVector],
start : Int,
end : Int,
period : Int,
policy : AlignmentPolicy,
) -> (Array[AlignedVector], AlignmentReport) {
align_timed_vectors(samples, regular_timestamps(start, end, period), policy)
}
///|
/// Estimate a constant lag by minimizing nearest-neighbor squared error.
pub fn estimate_stream_lag(
source : Array[TimedVector],
reference : Array[TimedVector],
candidate_lags : Array[Int],
policy : AlignmentPolicy,
) -> Int {
if candidate_lags.length() == 0 {
return 0
}
let mut best_lag = candidate_lags[0]
let mut best_error = 1.0e300
for lag in candidate_lags {
let mut error = 0.0
let mut count = 0
for sample in source {
let target = timed_vector_nearest(
reference,
sample.timestamp() + lag,
policy,
)
if target.is_present() &&
target.values().length() == sample.values().length() {
error = error +
vector_distance(sample.values(), target.values()) *
vector_distance(sample.values(), target.values())
count = count + 1
}
}
if count > 0 {
error = error / count.to_double()
}
if error < best_error {
best_error = error
best_lag = lag
}
}
best_lag
}
///|
/// Shift a stream by a measured lag.
pub fn shift_timed_vector_stream(
samples : Array[TimedVector],
lag : Int,
) -> Array[TimedVector] {
Array::makei(samples.length(), i => samples[i].shift(lag))
}
///|
/// Join two streams at matching timestamps.
pub fn join_timed_vectors(
left : Array[TimedVector],
right : Array[TimedVector],
policy : AlignmentPolicy,
) -> Array[(Int, Array[Double], Array[Double], Double)] {
let result = []
for sample in left {
let aligned = timed_vector_nearest(right, sample.timestamp(), policy)
if aligned.is_present() {
result.push(
(
sample.timestamp(),
sample.values(),
aligned.values(),
sample.quality().min(aligned.quality()),
),
)
}
}
result
}
///|
/// Compute the fraction of a stream that meets a quality threshold.
pub fn timed_vector_quality_fraction(
samples : Array[TimedVector],
threshold : Double,
) -> Double {
if samples.length() == 0 {
return 0.0
}
let mut count = 0
for sample in samples {
if sample.is_valid() && sample.quality() >= threshold {
count = count + 1
}
}
count.to_double() / samples.length().to_double()
}
///|
/// Compute the average vector over finite samples.
pub fn timed_vector_mean(samples : Array[TimedVector]) -> Array[Double] {
if samples.length() == 0 {
return []
}
let dimension = samples[0].dimension()
let sum = Array::make(dimension, 0.0)
let mut count = 0
for sample in samples {
if sample.is_valid() && sample.dimension() == dimension {
for i in 0.. Matrix {
let mean = timed_vector_mean(samples)
if mean.length() == 0 {
return Matrix::zeros(0, 0)
}
let covariance = Matrix::zeros(mean.length(), mean.length())
let mut count = 0
for sample in samples {
if sample.is_valid() && sample.dimension() == mean.length() {
let delta = vector_sub(sample.values(), mean)
for i in 0.. ignore
}
}
count = count + 1
}
}
if count > 1 {
covariance.scale(1.0 / (count - 1).to_double())
} else {
covariance
}
}
///|
/// Return an alignment report for a nearest-neighbor join.
pub fn join_alignment_report(
left : Array[TimedVector],
right : Array[TimedVector],
policy : AlignmentPolicy,
) -> AlignmentReport {
let timestamps = Array::makei(left.length(), i => left[i].timestamp())
let (_, report) = align_timed_vectors(right, timestamps, policy)
report
}
///|
/// Return the maximum absolute quality drop in a stream.
pub fn timed_vector_quality_drop(samples : Array[TimedVector]) -> Double {
if samples.length() < 2 {
return 0.0
}
let mut result = 0.0
for i in 1.. result {
result = drop
}
}
result
}
///|
/// Return a stream with quality attenuated by timestamp distance from a
/// reference time. This is useful when delayed samples should be retained
/// but contribute less to a fusion policy.
pub fn attenuate_timed_vector_quality(
samples : Array[TimedVector],
reference : Int,
time_constant : Double,
) -> Array[TimedVector] {
let constant = if time_constant <= 0.0 { 1.0 } else { time_constant }
Array::makei(samples.length(), i => {
let age = (samples[i].timestamp() - reference).abs().to_double()
let factor = 1.0 / (1.0 + age / constant)
samples[i].with_quality(samples[i].quality() * factor)
})
}
///|
/// Return a regular grid that covers all valid samples.
pub fn timed_vector_grid(
samples : Array[TimedVector],
period : Int,
) -> Array[Int] {
let filtered = filter_timed_vectors(samples, 0.0)
if filtered.length() == 0 {
return []
}
let sorted = sort_timed_vectors(filtered)
regular_timestamps(
sorted[0].timestamp(),
sorted[sorted.length() - 1].timestamp(),
period,
)
}
///|
/// Compute an alignment score combining coverage, quality, and distance.
pub fn alignment_quality_score(
report : AlignmentReport,
average_quality : Double,
distance_limit : Int,
) -> Double {
let coverage = report.coverage()
let quality = average_quality.clamp(min=0.0, max=1.0)
let distance = if distance_limit <= 0 {
1.0
} else {
(1.0 - report.mean_distance() / distance_limit.to_double()).clamp(
min=0.0,
max=1.0,
)
}
coverage * quality * distance
}