///|
/// A dependency-free in-memory telemetry store for CLI tools, tests, and
/// embedded applications. It provides deduplication, bounded retention,
/// deterministic queries, and stable export without pretending to be a DB.
pub(all) struct TelemetryRecord {
record_id : String
subject_id : String
source_id : String
timestamp_seconds : Double
rr_ms : Double
heart_rate_bpm : Double
movement_g : Double
temperature_c : Double
signal_quality : Double
tags : Array[String]
} derive(FromJson, ToJson, Debug, Eq)
///|
pub(all) struct TelemetryStorePolicy {
maximum_records : Int
minimum_quality : Double
reject_duplicate_ids : Bool
reject_non_monotonic_subjects : Bool
retention_seconds : Double
} derive(FromJson, ToJson, Debug, Eq)
///|
pub fn TelemetryStorePolicy::default() -> TelemetryStorePolicy {
{
maximum_records: 100000,
minimum_quality: 0.0,
reject_duplicate_ids: true,
reject_non_monotonic_subjects: false,
retention_seconds: 0.0,
}
}
///|
pub(all) struct TelemetryStoreNotice {
code : String
record_id : String
accepted : Bool
message : String
} derive(FromJson, ToJson, Debug, Eq)
///|
pub(all) struct TelemetryStore {
mut records : Array[TelemetryRecord]
mut notices : Array[TelemetryStoreNotice]
policy : TelemetryStorePolicy
} derive(Debug)
///|
pub(all) struct TelemetryQuery {
subject_id : String?
source_id : String?
start_seconds : Double?
end_seconds : Double?
minimum_quality : Double
limit : Int
descending : Bool
} derive(FromJson, ToJson, Debug, Eq)
///|
pub(all) struct TelemetrySummary {
record_count : Int
subject_count : Int
source_count : Int
first_timestamp : Double
last_timestamp : Double
duration_seconds : Double
mean_rr_ms : Double
mean_heart_rate_bpm : Double
mean_quality : Double
low_quality_count : Int
duplicate_notice_count : Int
} derive(FromJson, ToJson, Debug, Eq)
///|
pub(all) struct TelemetryRetentionResult {
before_count : Int
after_count : Int
removed_count : Int
cutoff_seconds : Double
} derive(FromJson, ToJson, Debug, Eq)
///|
fn telemetry_bound(value : Double, low : Double, high : Double) -> Double {
if value.is_nan() || value.is_inf() {
low
} else {
value.clamp(min=low, max=high)
}
}
///|
pub fn make_telemetry_record(
record_id : String,
subject_id : String,
source_id : String,
timestamp_seconds : Double,
rr_ms : Double,
heart_rate_bpm : Double,
movement_g : Double,
temperature_c : Double,
signal_quality : Double,
tags : Array[String],
) -> TelemetryRecord {
{
record_id,
subject_id,
source_id,
timestamp_seconds: timestamp_seconds.max(0.0),
rr_ms: telemetry_bound(rr_ms, 0.0, 3000.0),
heart_rate_bpm: telemetry_bound(heart_rate_bpm, 0.0, 260.0),
movement_g: telemetry_bound(movement_g, 0.0, 100.0),
temperature_c: telemetry_bound(temperature_c, -50.0, 100.0),
signal_quality: telemetry_bound(signal_quality, 0.0, 1.0),
tags,
}
}
///|
pub fn telemetry_record_is_valid(record : TelemetryRecord) -> Bool {
record.record_id.length() > 0 &&
record.subject_id.length() > 0 &&
record.source_id.length() > 0 &&
record.timestamp_seconds >= 0.0 &&
record.rr_ms >= 0.0 &&
record.heart_rate_bpm >= 0.0 &&
record.signal_quality >= 0.0 &&
record.signal_quality <= 1.0
}
///|
pub fn telemetry_record_key(record : TelemetryRecord) -> String {
"\{record.subject_id}|\{record.source_id}|\{record.record_id}"
}
///|
pub fn telemetry_record_row(record : TelemetryRecord) -> Array[String] {
[
record.record_id,
record.subject_id,
record.source_id,
record.timestamp_seconds.to_string(),
record.rr_ms.to_string(),
record.heart_rate_bpm.to_string(),
record.movement_g.to_string(),
record.temperature_c.to_string(),
record.signal_quality.to_string(),
record.tags.join("|"),
]
}
///|
pub fn TelemetryStore::new(policy : TelemetryStorePolicy) -> TelemetryStore {
{ records: [], notices: [], policy }
}
///|
fn telemetry_notice(
store : TelemetryStore,
record : TelemetryRecord,
accepted : Bool,
code : String,
message : String,
) -> Unit {
store.notices.push({ code, record_id: record.record_id, accepted, message })
}
///|
fn telemetry_has_id(store : TelemetryStore, record : TelemetryRecord) -> Bool {
let key = telemetry_record_key(record)
store.records.any(existing => telemetry_record_key(existing) == key)
}
///|
fn telemetry_subject_last(
store : TelemetryStore,
subject_id : String,
) -> Double? {
let mut found : Double? = None
for record in store.records {
if record.subject_id == subject_id {
match found {
Some(value) =>
if record.timestamp_seconds > value {
found = Some(record.timestamp_seconds)
}
None => found = Some(record.timestamp_seconds)
}
}
}
found
}
///|
pub fn TelemetryStore::insert(
self : TelemetryStore,
record : TelemetryRecord,
) -> Bool {
if !telemetry_record_is_valid(record) {
telemetry_notice(
self, record, false, "invalid_record", "record failed schema validation",
)
false
} else if record.signal_quality < self.policy.minimum_quality {
telemetry_notice(
self, record, false, "low_quality", "record is below the store quality floor",
)
false
} else if self.policy.reject_duplicate_ids && telemetry_has_id(self, record) {
telemetry_notice(
self, record, false, "duplicate_id", "record key already exists",
)
false
} else {
match telemetry_subject_last(self, record.subject_id) {
Some(last) =>
if self.policy.reject_non_monotonic_subjects &&
record.timestamp_seconds < last {
telemetry_notice(
self, record, false, "non_monotonic", "timestamp precedes the latest subject record",
)
false
} else {
self.records.push(record)
true
}
None => {
self.records.push(record)
true
}
}
}
}
///|
pub fn TelemetryStore::insert_many(
self : TelemetryStore,
records : Array[TelemetryRecord],
) -> Int {
let mut accepted = 0
for record in records {
if self.insert(record) {
accepted += 1
}
}
accepted
}
///|
pub fn TelemetryStore::count(self : TelemetryStore) -> Int {
self.records.length()
}
///|
pub fn TelemetryStore::notice_count(self : TelemetryStore) -> Int {
self.notices.length()
}
///|
pub fn TelemetryStore::records(self : TelemetryStore) -> Array[TelemetryRecord] {
let result = []
for record in self.records {
result.push(record)
}
result
}
///|
pub fn TelemetryStore::notices(
self : TelemetryStore,
) -> Array[TelemetryStoreNotice] {
let result = []
for notice in self.notices {
result.push(notice)
}
result
}
///|
fn telemetry_match_optional(value : String, wanted : String?) -> Bool {
match wanted {
Some(expected) => value == expected
None => true
}
}
///|
fn telemetry_match_time(value : Double, start : Double?, end : Double?) -> Bool {
let after_start = match start {
Some(bound) => value >= bound
None => true
}
let before_end = match end {
Some(bound) => value <= bound
None => true
}
after_start && before_end
}
///|
fn telemetry_copy(records : Array[TelemetryRecord]) -> Array[TelemetryRecord] {
let result = []
for record in records {
result.push(record)
}
result
}
///|
pub fn TelemetryStore::query(
self : TelemetryStore,
query : TelemetryQuery,
) -> Array[TelemetryRecord] {
let mut matched = []
for record in self.records {
if telemetry_match_optional(record.subject_id, query.subject_id) &&
telemetry_match_optional(record.source_id, query.source_id) &&
telemetry_match_time(
record.timestamp_seconds,
query.start_seconds,
query.end_seconds,
) &&
record.signal_quality >= query.minimum_quality {
matched.push(record)
}
}
matched.sort_by((left, right) => {
if left.timestamp_seconds < right.timestamp_seconds {
-1
} else if left.timestamp_seconds > right.timestamp_seconds {
1
} else {
0
}
})
if query.descending {
let reversed = []
for i in 0..= 0 && matched.length() > query.limit {
matched.truncate(query.limit)
}
matched
}
///|
pub fn telemetry_query_all() -> TelemetryQuery {
{
subject_id: None,
source_id: None,
start_seconds: None,
end_seconds: None,
minimum_quality: 0.0,
limit: -1,
descending: false,
}
}
///|
pub fn telemetry_query_subject(
subject_id : String,
limit : Int,
) -> TelemetryQuery {
{
subject_id: Some(subject_id),
source_id: None,
start_seconds: None,
end_seconds: None,
minimum_quality: 0.0,
limit,
descending: false,
}
}
///|
pub fn telemetry_query_window(
subject_id : String,
start_seconds : Double,
end_seconds : Double,
) -> TelemetryQuery {
{
subject_id: Some(subject_id),
source_id: None,
start_seconds: Some(start_seconds.min(end_seconds)),
end_seconds: Some(start_seconds.max(end_seconds)),
minimum_quality: 0.0,
limit: -1,
descending: false,
}
}
///|
fn telemetry_unique_strings(values : Array[String]) -> Array[String] {
let result = []
for value in values {
if !result.any(existing => existing == value) {
result.push(value)
}
}
result
}
///|
fn telemetry_min_timestamp(values : Array[Double]) -> Double {
if values.length() == 0 {
0.0
} else {
let mut result = values[0]
for value in values {
if value < result {
result = value
}
}
result
}
}
///|
fn telemetry_max_timestamp(values : Array[Double]) -> Double {
if values.length() == 0 {
0.0
} else {
let mut result = values[0]
for value in values {
if value > result {
result = value
}
}
result
}
}
///|
pub fn telemetry_summary(records : Array[TelemetryRecord]) -> TelemetrySummary {
let subjects = telemetry_unique_strings(
records.map(record => record.subject_id),
)
let sources = telemetry_unique_strings(
records.map(record => record.source_id),
)
let timestamps = records.map(record => record.timestamp_seconds)
let rr = records
.filter(record => record.rr_ms > 0.0)
.map(record => record.rr_ms)
let hr = records
.filter(record => record.heart_rate_bpm > 0.0)
.map(record => record.heart_rate_bpm)
let quality = records.map(record => record.signal_quality)
let first_timestamp = if timestamps.length() == 0 {
0.0
} else {
telemetry_min_timestamp(timestamps)
}
let last_timestamp = if timestamps.length() == 0 {
0.0
} else {
telemetry_max_timestamp(timestamps)
}
{
record_count: records.length(),
subject_count: subjects.length(),
source_count: sources.length(),
first_timestamp,
last_timestamp,
duration_seconds: last_timestamp - first_timestamp,
mean_rr_ms: mean_value(rr),
mean_heart_rate_bpm: mean_value(hr),
mean_quality: mean_value(quality),
low_quality_count: records
.filter(record => record.signal_quality < 0.70)
.length(),
duplicate_notice_count: 0,
}
}
///|
pub fn telemetry_store_summary(store : TelemetryStore) -> TelemetrySummary {
let summary = telemetry_summary(store.records)
{
record_count: summary.record_count,
subject_count: summary.subject_count,
source_count: summary.source_count,
first_timestamp: summary.first_timestamp,
last_timestamp: summary.last_timestamp,
duration_seconds: summary.duration_seconds,
mean_rr_ms: summary.mean_rr_ms,
mean_heart_rate_bpm: summary.mean_heart_rate_bpm,
mean_quality: summary.mean_quality,
low_quality_count: summary.low_quality_count,
duplicate_notice_count: store.notices
.filter(notice => notice.code == "duplicate_id")
.length(),
}
}
///|
pub fn TelemetryStore::retain_after(
self : TelemetryStore,
cutoff_seconds : Double,
) -> TelemetryRetentionResult {
let before = self.records.length()
let kept = self.records.filter(record => {
record.timestamp_seconds >= cutoff_seconds
})
self.records = kept
{
before_count: before,
after_count: kept.length(),
removed_count: before - kept.length(),
cutoff_seconds,
}
}
///|
pub fn TelemetryStore::enforce_limits(self : TelemetryStore) -> Int {
if self.policy.maximum_records >= 0 &&
self.records.length() > self.policy.maximum_records {
self.records.sort_by((left, right) => {
if left.timestamp_seconds < right.timestamp_seconds {
-1
} else if left.timestamp_seconds > right.timestamp_seconds {
1
} else {
0
}
})
let remove_count = self.records.length() - self.policy.maximum_records
let kept = []
for i in remove_count.. TelemetryRetentionResult {
let cutoff = if self.policy.retention_seconds <= 0.0 {
0.0
} else {
(newest_timestamp - self.policy.retention_seconds).max(0.0)
}
let result = self.retain_after(cutoff)
let _ = self.enforce_limits()
result
}
///|
pub fn telemetry_export_csv(records : Array[TelemetryRecord]) -> String {
let grid = [
[
"record_id", "subject_id", "source_id", "timestamp_seconds", "rr_ms", "heart_rate_bpm",
"movement_g", "temperature_c", "signal_quality", "tags",
],
]
for record in records {
grid.push(telemetry_record_row(record))
}
to_csv(grid)
}
///|
pub fn TelemetryStore::export_csv(self : TelemetryStore) -> String {
telemetry_export_csv(self.records)
}
///|
pub fn telemetry_notices_csv(notices : Array[TelemetryStoreNotice]) -> String {
let grid = [["code", "record_id", "accepted", "message"]]
for notice in notices {
grid.push([
notice.code,
notice.record_id,
notice.accepted.to_string(),
notice.message,
])
}
to_csv(grid)
}
///|
pub fn telemetry_store_feature_vector(store : TelemetryStore) -> Array[Double] {
let summary = telemetry_store_summary(store)
[
summary.record_count.to_double(),
summary.subject_count.to_double(),
summary.source_count.to_double(),
summary.duration_seconds,
summary.mean_rr_ms,
summary.mean_heart_rate_bpm,
summary.mean_quality,
summary.low_quality_count.to_double(),
summary.duplicate_notice_count.to_double(),
]
}
///|
pub fn telemetry_store_is_healthy(store : TelemetryStore) -> Bool {
let summary = telemetry_store_summary(store)
summary.record_count > 0 &&
summary.mean_quality >= store.policy.minimum_quality &&
summary.low_quality_count.to_double() <=
summary.record_count.to_double() * 0.30
}
///|
pub fn telemetry_records_for_source(
records : Array[TelemetryRecord],
source_id : String,
) -> Array[TelemetryRecord] {
records.filter(record => record.source_id == source_id)
}
///|
pub fn telemetry_records_for_subject(
records : Array[TelemetryRecord],
subject_id : String,
) -> Array[TelemetryRecord] {
records.filter(record => record.subject_id == subject_id)
}
///|
pub fn telemetry_gap_count(
records : Array[TelemetryRecord],
maximum_gap_seconds : Double,
) -> Int {
let ordered = telemetry_copy(records)
ordered.sort_by((left, right) => {
if left.timestamp_seconds < right.timestamp_seconds {
-1
} else if left.timestamp_seconds > right.timestamp_seconds {
1
} else {
0
}
})
let mut gaps = 0
for i in 1..
maximum_gap_seconds {
gaps += 1
}
}
gaps
}
///|
pub fn telemetry_quality_histogram(
records : Array[TelemetryRecord],
buckets : Int,
) -> Array[Int] {
let size = buckets.clamp(min=1, max=100)
let result = Array::make(size, 0)
for record in records {
let index = (record.signal_quality.clamp(min=0.0, max=0.999999) *
size.to_double()).to_int()
result[index] += 1
}
result
}
///|
pub fn telemetry_latest_by_subject(
records : Array[TelemetryRecord],
) -> Array[TelemetryRecord] {
let subjects = telemetry_unique_strings(
records.map(record => record.subject_id),
)
let result = []
for subject in subjects {
let matches = telemetry_records_for_subject(records, subject)
if matches.length() > 0 {
let mut latest = matches[0]
for record in matches {
if record.timestamp_seconds > latest.timestamp_seconds {
latest = record
}
}
result.push(latest)
}
}
result
}