///|
/// Full output of the standardization stage.
pub(all) struct NormalizationResult {
readings : Array[Reading]
gaps : Array[SamplingGap]
diagnostics : Array[Diagnostic]
} derive(Debug, Eq)
///|
/// Compare readings by sensor id, timestamp and source line.
fn compare_readings(a : Reading, b : Reading) -> Int {
let sensor_order = a.sensor_id.compare(b.sensor_id)
if sensor_order != 0 {
return sensor_order
}
if a.timestamp < b.timestamp {
return -1
}
if a.timestamp > b.timestamp {
return 1
}
match (a.source_line, b.source_line) {
(Some(left), Some(right)) => left - right
(Some(_), None) => -1
(None, Some(_)) => 1
(None, None) => 0
}
}
///|
/// Return an ordered copy without mutating caller storage.
pub fn sort_readings(readings : Array[Reading]) -> Array[Reading] {
let ordered = copy_array(readings)
ordered.sort_by(compare_readings)
ordered
}
///|
/// Average optional numeric values when both are present.
fn average_optional(a : Double?, b : Double?) -> Double? {
match (a, b) {
(Some(left), Some(right)) => Some((left + right) / 2.0)
(Some(left), None) => Some(left)
(None, Some(right)) => Some(right)
(None, None) => None
}
}
///|
/// Merge two readings representing the same sensor and instant.
fn merge_duplicate_pair(
first : Reading,
second : Reading,
average : Bool,
) -> Reading {
if average {
{
timestamp: first.timestamp,
sensor_id: first.sensor_id,
temperature_c: (first.temperature_c + second.temperature_c) / 2.0,
humidity_percent: average_optional(
first.humidity_percent,
second.humidity_percent,
),
battery_percent: average_optional(
first.battery_percent,
second.battery_percent,
),
status: if first.status == second.status {
first.status
} else {
"merged"
},
origin: first.origin,
flags: copy_array(first.flags),
source_line: first.source_line,
}.add_flag(DuplicateTimestamp)
} else {
second.add_flag(DuplicateTimestamp)
}
}
///|
/// Collapse duplicate sensor and timestamp keys.
pub fn merge_duplicate_readings(
ordered : Array[Reading],
average : Bool,
) -> ReadingBatch {
let output : Array[Reading] = []
let diagnostics : Array[Diagnostic] = []
for sample in ordered {
if output.length() == 0 {
output.push(sample)
continue
}
let previous = output[output.length() - 1]
if previous.sensor_id == sample.sensor_id &&
previous.timestamp == sample.timestamp {
output[output.length() - 1] = merge_duplicate_pair(
previous, sample, average,
)
diagnostics.push(
warning_diagnostic(
"normalize.duplicate_timestamp", "duplicate timestamp was merged",
).for_sensor(sample.sensor_id),
)
} else {
output.push(sample)
}
}
{ readings: output, diagnostics }
}
///|
/// Apply calibration and physical plausibility checks.
pub fn calibrate_and_validate(
readings : Array[Reading],
config : AnalysisConfig,
) -> ReadingBatch {
let output : Array[Reading] = []
let diagnostics : Array[Diagnostic] = []
for sample in readings {
let profile = config.profile_for(sample.sensor_id)
let calibrated_temperature = sample.temperature_c +
profile.calibration_offset_c
let mut next = { ..sample, temperature_c: calibrated_temperature }
if profile.calibration_offset_c != 0.0 {
next = next.add_flag(Calibrated)
}
if calibrated_temperature < profile.physical_min_c ||
calibrated_temperature > profile.physical_max_c {
next = next.add_flag(OutOfPhysicalRange)
diagnostics.push(
warning_diagnostic(
"normalize.physical_range", "temperature is outside the configured physical sensor range",
).for_sensor(sample.sensor_id),
)
}
match next.humidity_percent {
Some(value) =>
if value < 0.0 || value > 100.0 {
diagnostics.push(
warning_diagnostic(
"normalize.humidity_range", "humidity is outside 0–100 percent",
).for_sensor(sample.sensor_id),
)
}
None => ()
}
match next.battery_percent {
Some(value) =>
if value < 0.0 || value > 100.0 {
diagnostics.push(
warning_diagnostic(
"normalize.battery_range", "battery is outside 0–100 percent",
).for_sensor(sample.sensor_id),
)
}
None => ()
}
output.push(next)
}
{ readings: output, diagnostics }
}
///|
/// Estimate the number of samples missing from one interval.
fn estimated_missing_count(duration : Int64, expected : Int64) -> Int {
if expected <= 0L || duration <= expected {
0
} else {
let intervals = duration / expected
if intervals <= 1L {
0
} else {
(intervals - 1L).to_int()
}
}
}
///|
/// Find sampling gaps and tag the first sample after each gap.
pub fn detect_sampling_gaps(
readings : Array[Reading],
config : AnalysisConfig,
) -> NormalizationResult {
let output = copy_array(readings)
let gaps : Array[SamplingGap] = []
let diagnostics : Array[Diagnostic] = []
for index = 1; index < output.length(); index = index + 1 {
let previous = output[index - 1]
let current = output[index]
if previous.sensor_id != current.sensor_id {
continue
}
let profile = config.profile_for(current.sensor_id)
let duration = current.timestamp - previous.timestamp
let threshold = if config.policy.maximum_gap_seconds >
profile.expected_interval_seconds {
config.policy.maximum_gap_seconds
} else {
profile.expected_interval_seconds * 2L
}
if duration > threshold {
gaps.push({
sensor_id: current.sensor_id,
previous_timestamp: previous.timestamp,
next_timestamp: current.timestamp,
duration_seconds: duration,
estimated_missing_samples: estimated_missing_count(
duration,
profile.expected_interval_seconds,
),
})
output[index] = current.add_flag(GapBefore)
diagnostics.push(
warning_diagnostic(
"normalize.sampling_gap", "sampling interval exceeded the configured maximum",
).for_sensor(current.sensor_id),
)
}
}
{ readings: output, gaps, diagnostics }
}
///|
/// Linear interpolation helper.
fn interpolate_value(
start : Double,
end : Double,
numerator : Int64,
denominator : Int64,
) -> Double {
start + (end - start) * numerator.to_double() / denominator.to_double()
}
///|
/// Interpolate optional values only when both endpoints are present.
fn interpolate_optional(
start : Double?,
end : Double?,
numerator : Int64,
denominator : Int64,
) -> Double? {
match (start, end) {
(Some(left), Some(right)) =>
Some(interpolate_value(left, right, numerator, denominator))
_ => None
}
}
///|
/// Insert expected samples in short gaps. Long gaps remain explicit.
pub fn interpolate_short_gaps(
readings : Array[Reading],
config : AnalysisConfig,
) -> ReadingBatch {
if !config.interpolate_short_gaps || readings.length() < 2 {
return { readings: copy_array(readings), diagnostics: [] }
}
let output : Array[Reading] = []
let diagnostics : Array[Diagnostic] = []
output.push(readings[0])
for index = 1; index < readings.length(); index = index + 1 {
let previous = readings[index - 1]
let current = readings[index]
if previous.sensor_id == current.sensor_id {
let profile = config.profile_for(current.sensor_id)
let duration = current.timestamp - previous.timestamp
let interval = profile.expected_interval_seconds
if interval > 0L &&
duration > interval &&
duration <= config.interpolation_limit_seconds {
let mut offset = interval
while offset < duration {
output.push({
timestamp: previous.timestamp + offset,
sensor_id: previous.sensor_id,
temperature_c: interpolate_value(
previous.temperature_c,
current.temperature_c,
offset,
duration,
),
humidity_percent: interpolate_optional(
previous.humidity_percent,
current.humidity_percent,
offset,
duration,
),
battery_percent: interpolate_optional(
previous.battery_percent,
current.battery_percent,
offset,
duration,
),
status: "interpolated",
origin: Interpolated,
flags: [],
source_line: None,
})
offset = offset + interval
}
diagnostics.push(
info_diagnostic(
"normalize.interpolated_gap", "short sampling gap was linearly interpolated",
).for_sensor(current.sensor_id),
)
}
}
output.push(current)
}
{ readings: output, diagnostics }
}
///|
/// Run the deterministic standardization pipeline.
pub fn normalize_readings(
readings : Array[Reading],
config : AnalysisConfig,
) -> NormalizationResult {
let diagnostics = validate_analysis_config(config)
let ordered = sort_readings(readings)
let duplicates = merge_duplicate_readings(
ordered,
config.merge_duplicate_average,
)
diagnostics.append(duplicates.diagnostics)
let calibrated = calibrate_and_validate(duplicates.readings, config)
diagnostics.append(calibrated.diagnostics)
let interpolated = interpolate_short_gaps(calibrated.readings, config)
diagnostics.append(interpolated.diagnostics)
let gaps = detect_sampling_gaps(interpolated.readings, config)
diagnostics.append(gaps.diagnostics)
{ readings: gaps.readings, gaps: gaps.gaps, diagnostics }
}
///|
/// Return all readings for one sensor.
pub fn readings_for_sensor(
readings : Array[Reading],
sensor_id : String,
) -> Array[Reading] {
let selected : Array[Reading] = []
for sample in readings {
if sample.sensor_id == sensor_id {
selected.push(sample)
}
}
selected
}
///|
/// Return stable, sorted unique sensor identifiers.
pub fn unique_sensor_ids(readings : Array[Reading]) -> Array[String] {
let ids : Array[String] = []
for sample in readings {
let mut found = false
for existing in ids {
if existing == sample.sensor_id {
found = true
}
}
if !found {
ids.push(sample.sensor_id)
}
}
ids.sort()
ids
}
///|
/// Determine total observation span per sensor.
pub fn observation_span(readings : Array[Reading], sensor_id : String) -> Int64 {
let selected = readings_for_sensor(readings, sensor_id)
if selected.length() < 2 {
0L
} else {
let ordered = sort_readings(selected)
ordered[ordered.length() - 1].timestamp - ordered[0].timestamp
}
}
///|
/// Count readings by origin.
pub fn count_origin(readings : Array[Reading], origin : SampleOrigin) -> Int {
let mut count = 0
for sample in readings {
if sample.origin == origin {
count = count + 1
}
}
count
}