///|
/// CSV protocol adapters for common wearable exports.
/// The adapter layer separates source-specific column names and units from
/// the source-independent WearableSample and ingestion pipeline.
pub(all) enum WearableProtocolDialect {
ProtocolGeneric
ProtocolPolar
ProtocolGarmin
ProtocolAppleHealth
ProtocolFitExport
} derive(FromJson, ToJson, Debug, Eq)
///|
pub(all) enum ProtocolUnit {
ProtocolMilliseconds
ProtocolSeconds
ProtocolBeatsPerMinute
ProtocolKilometers
ProtocolMeters
ProtocolG
ProtocolCelsius
ProtocolPercent
ProtocolUnknown
} derive(FromJson, ToJson, Debug, Eq)
///|
pub(all) struct ProtocolColumnMapping {
timestamp_index : Int
rr_index : Int
heart_rate_index : Int
quality_index : Int
movement_index : Int
temperature_index : Int
timestamp_unit : ProtocolUnit
rr_unit : ProtocolUnit
quality_unit : ProtocolUnit
} derive(FromJson, ToJson, Debug, Eq)
///|
pub(all) struct ProtocolAdapterConfig {
dialect : WearableProtocolDialect
delimiter : String
has_header : Bool
default_quality : Double
source_id : String
start_timestamp_seconds : Double
infer_heart_rate : Bool
infer_quality : Bool
} derive(FromJson, ToJson, Debug, Eq)
///|
pub fn ProtocolAdapterConfig::default() -> ProtocolAdapterConfig {
{
dialect: ProtocolGeneric,
delimiter: ",",
has_header: true,
default_quality: 1.0,
source_id: "csv",
start_timestamp_seconds: 0.0,
infer_heart_rate: true,
infer_quality: true,
}
}
///|
pub(all) struct ProtocolAdapterNotice {
row_index : Int
code : String
severity : String
message : String
} derive(FromJson, ToJson, Debug, Eq)
///|
pub(all) struct ProtocolAdapterResult {
dialect : WearableProtocolDialect
mapping : ProtocolColumnMapping
samples : Array[WearableSample]
notices : Array[ProtocolAdapterNotice]
accepted_count : Int
rejected_count : Int
header : Array[String]
source_id : String
} derive(FromJson, ToJson, Debug, Eq)
///|
fn protocol_dialect_name(value : WearableProtocolDialect) -> String {
match value {
ProtocolGeneric => "generic"
ProtocolPolar => "polar"
ProtocolGarmin => "garmin"
ProtocolAppleHealth => "apple_health"
ProtocolFitExport => "fit_export"
}
}
///|
pub fn protocol_unit_name(value : ProtocolUnit) -> String {
match value {
ProtocolMilliseconds => "milliseconds"
ProtocolSeconds => "seconds"
ProtocolBeatsPerMinute => "bpm"
ProtocolKilometers => "kilometers"
ProtocolMeters => "meters"
ProtocolG => "g"
ProtocolCelsius => "celsius"
ProtocolPercent => "percent"
ProtocolUnknown => "unknown"
}
}
///|
fn protocol_clean(value : String) -> String {
trim(value).to_owned()
}
///|
fn protocol_number(value : String, fallback : Double) -> Double {
@strconv.from_str(trim(value)) catch {
_ => fallback
}
}
///|
fn protocol_index(headers : Array[String], names : Array[String]) -> Int {
for i in 0.. ProtocolColumnMapping {
let timestamp_names = match dialect {
ProtocolPolar => ["timestamp", "time", "sample_time"]
ProtocolGarmin => ["timestamp", "time", "date_time"]
ProtocolAppleHealth => ["startdate", "date", "timestamp"]
ProtocolFitExport => ["timestamp", "time", "elapsed_time"]
ProtocolGeneric => ["timestamp", "time", "seconds", "time_seconds"]
}
let rr_names = ["rr_ms", "rr", "rr_interval", "ibi", "interbeat_interval"]
let hr_names = ["heart_rate", "heart_rate_bpm", "hr", "pulse"]
let quality_names = [
"quality", "signal_quality", "confidence", "confidence_percent",
]
let movement_names = ["movement", "movement_g", "acceleration", "accel"]
let temperature_names = ["temperature", "temperature_c", "skin_temperature"]
{
timestamp_index: protocol_index(headers, timestamp_names),
rr_index: protocol_index(headers, rr_names),
heart_rate_index: protocol_index(headers, hr_names),
quality_index: protocol_index(headers, quality_names),
movement_index: protocol_index(headers, movement_names),
temperature_index: protocol_index(headers, temperature_names),
timestamp_unit: ProtocolSeconds,
rr_unit: if dialect == ProtocolPolar {
ProtocolMilliseconds
} else {
ProtocolMilliseconds
},
quality_unit: if dialect == ProtocolAppleHealth {
ProtocolPercent
} else {
ProtocolPercent
},
}
}
///|
pub fn protocol_mapping_is_usable(mapping : ProtocolColumnMapping) -> Bool {
mapping.timestamp_index >= 0 && mapping.rr_index >= 0
}
///|
fn protocol_row_value(row : Array[String], index : Int) -> String {
if index < 0 || index >= row.length() {
""
} else {
trim(row[index]).to_owned()
}
}
///|
fn protocol_unit_value(value : Double, unit : ProtocolUnit) -> Double {
match unit {
ProtocolMilliseconds => value
ProtocolSeconds => value * 1000.0
ProtocolPercent => value / 100.0
_ => value
}
}
///|
fn protocol_timestamp_value(
value : Double,
unit : ProtocolUnit,
fallback : Double,
) -> Double {
let adjusted = match unit {
ProtocolMilliseconds => value / 1000.0
ProtocolSeconds => value
_ => value
}
if adjusted.is_nan() || adjusted.is_inf() || adjusted < 0.0 {
fallback
} else {
adjusted
}
}
///|
fn protocol_notice(
notices : Array[ProtocolAdapterNotice],
row_index : Int,
code : String,
severity : String,
message : String,
) -> Unit {
notices.push({ row_index, code, severity, message })
}
///|
pub fn protocol_normalize_quality(
value : Double,
unit : ProtocolUnit,
fallback : Double,
) -> Double {
let normalized = protocol_unit_value(value, unit)
if normalized.is_nan() || normalized.is_inf() {
fallback
} else if unit == ProtocolPercent {
normalized.clamp(min=0.0, max=1.0)
} else {
normalized.clamp(min=0.0, max=1.0)
}
}
///|
pub fn protocol_parse_csv(
content : String,
config : ProtocolAdapterConfig,
) -> ProtocolAdapterResult {
let grid = parse_csv(content)
let header = if config.has_header && grid.length() > 0 { grid[0] } else { [] }
let effective_header = if header.length() > 0 {
header
} else {
["timestamp", "rr_ms", "heart_rate_bpm", "quality"]
}
let mapping = protocol_mapping_for_headers(config.dialect, effective_header)
let notices = []
let samples = []
let start_row = if config.has_header { 1 } else { 0 }
let mut rejected = 0
if !protocol_mapping_is_usable(mapping) {
protocol_notice(
notices, -1, "missing_columns", "error", "timestamp and RR columns are required",
)
return {
dialect: config.dialect,
mapping,
samples,
notices,
accepted_count: 0,
rejected_count: grid.length(),
header: effective_header,
source_id: config.source_id,
}
}
for row_index in start_row.. 0.0 {
if mapping.heart_rate_index >= 0 {
hr_raw
} else {
60000.0 / rr_ms
}
} else if config.infer_heart_rate {
60000.0 / rr_ms
} else {
0.0
}
let quality_text = protocol_row_value(row, mapping.quality_index)
let quality_raw = protocol_number(quality_text, config.default_quality)
let quality = if mapping.quality_index >= 0 {
protocol_normalize_quality(
quality_raw,
mapping.quality_unit,
config.default_quality,
)
} else if config.infer_quality {
if rr_ms >= 300.0 && rr_ms <= 2000.0 {
config.default_quality
} else {
0.40
}
} else {
config.default_quality
}
let movement = protocol_number(
protocol_row_value(row, mapping.movement_index),
0.0,
)
let temperature = protocol_number(
protocol_row_value(row, mapping.temperature_index),
0.0,
)
let base_sample = make_wearable_sample(
timestamp,
rr_ms,
heart_rate,
quality,
config.source_id,
row_index - start_row,
)
samples.push(
wearable_sample_with_context(base_sample, movement, temperature),
)
}
{
dialect: config.dialect,
mapping,
samples,
notices,
accepted_count: samples.length(),
rejected_count: rejected,
header: effective_header,
source_id: config.source_id,
}
}
///|
pub fn protocol_result_is_usable(result : ProtocolAdapterResult) -> Bool {
result.accepted_count > 0 &&
protocol_mapping_is_usable(result.mapping) &&
result.samples.all(sample => {
wearable_sample_is_valid(sample, WearableIngestConfig::default())
})
}
///|
pub fn protocol_result_feature_vector(
result : ProtocolAdapterResult,
) -> Array[Double] {
[
result.accepted_count.to_double(),
result.rejected_count.to_double(),
result.notices.length().to_double(),
mean_value(result.samples.map(sample => sample.signal_quality)),
result.mapping.timestamp_index.to_double(),
result.mapping.rr_index.to_double(),
result.mapping.heart_rate_index.to_double(),
]
}
///|
pub fn protocol_result_csv(result : ProtocolAdapterResult) -> String {
wearable_ingest_csv(
ingest_wearable_samples(result.samples, WearableIngestConfig::default()),
)
}
///|
pub fn protocol_notices_csv(result : ProtocolAdapterResult) -> String {
let grid = [["row_index", "code", "severity", "message"]]
for notice in result.notices {
grid.push([
notice.row_index.to_string(),
notice.code,
notice.severity,
notice.message,
])
}
to_csv(grid)
}
///|
pub fn protocol_mapping_csv(mapping : ProtocolColumnMapping) -> String {
to_csv([
[
"timestamp_index", "rr_index", "heart_rate_index", "quality_index", "movement_index",
"temperature_index", "timestamp_unit", "rr_unit", "quality_unit",
],
[
mapping.timestamp_index.to_string(),
mapping.rr_index.to_string(),
mapping.heart_rate_index.to_string(),
mapping.quality_index.to_string(),
mapping.movement_index.to_string(),
mapping.temperature_index.to_string(),
protocol_unit_name(mapping.timestamp_unit),
protocol_unit_name(mapping.rr_unit),
protocol_unit_name(mapping.quality_unit),
],
])
}
///|
pub fn protocol_dialect_from_name(name : String) -> WearableProtocolDialect {
match protocol_clean(name) {
"polar" => ProtocolPolar
"garmin" => ProtocolGarmin
"apple" | "apple_health" => ProtocolAppleHealth
"fit" | "fit_export" => ProtocolFitExport
_ => ProtocolGeneric
}
}
///|
pub fn protocol_config_for_name(
name : String,
source_id : String,
) -> ProtocolAdapterConfig {
let default = ProtocolAdapterConfig::default()
{
dialect: protocol_dialect_from_name(name),
delimiter: default.delimiter,
has_header: default.has_header,
default_quality: default.default_quality,
source_id,
start_timestamp_seconds: default.start_timestamp_seconds,
infer_heart_rate: default.infer_heart_rate,
infer_quality: default.infer_quality,
}
}
///|
pub fn protocol_samples_to_csv(samples : Array[WearableSample]) -> String {
wearable_ingest_csv(
ingest_wearable_samples(samples, WearableIngestConfig::default()),
)
}
///|
pub fn protocol_roundtrip_quality(
content : String,
config : ProtocolAdapterConfig,
) -> Double {
let result = protocol_parse_csv(content, config)
if result.samples.length() == 0 {
0.0
} else {
mean_value(result.samples.map(sample => sample.signal_quality))
}
}
///|
pub fn protocol_rejection_ratio(result : ProtocolAdapterResult) -> Double {
let total = result.accepted_count + result.rejected_count
if total == 0 {
0.0
} else {
result.rejected_count.to_double() / total.to_double()
}
}
///|
pub fn protocol_notice_count_for(
result : ProtocolAdapterResult,
code : String,
) -> Int {
result.notices.filter(notice => notice.code == code).length()
}
///|
pub fn protocol_has_errors(result : ProtocolAdapterResult) -> Bool {
result.notices.any(notice => notice.severity == "error")
}
///|
pub fn protocol_result_summary(result : ProtocolAdapterResult) -> String {
"\{protocol_dialect_name(result.dialect)}: \{result.accepted_count.to_string()} accepted, \{result.rejected_count.to_string()} rejected"
}