use chrono::DateTime;
use serde::de::DeserializeOwned;
use serde::{Deserialize, Serialize};
use serde_json::Value;
use crate::config::Site;
const PRED_STEP_MS: i64 = 5 * 60_000;
const MAX_RESPONSE_BYTES: u64 = 8 * 1024 * 1024;
const MAX_FUTURE_SKEW_MS: i64 = 5 * 60_000;
const MIN_PHYSIOLOGICAL_SGV: f64 = 39.0;
pub type Result<T> = std::result::Result<T, NsError>;
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct Entry {
pub sgv: f64,
pub date: i64,
#[serde(default)]
pub direction: Option<String>,
}
impl Entry {
pub fn arrow(&self) -> &'static str {
match self.direction.as_deref() {
Some("DoubleUp") => "⇈",
Some("SingleUp") => "↑",
Some("FortyFiveUp") => "↗",
Some("Flat") => "→",
Some("FortyFiveDown") => "↘",
Some("SingleDown") => "↓",
Some("DoubleDown") => "⇊",
_ => "-",
}
}
}
#[derive(Clone)]
pub struct Client {
http: reqwest::Client,
base_url: String,
token: String,
}
impl Client {
pub fn for_site(site: &Site) -> Result<Self> {
let http = reqwest::Client::builder()
.user_agent(concat!("sugarrush/", env!("CARGO_PKG_VERSION")))
.timeout(std::time::Duration::from_secs(12))
.connect_timeout(std::time::Duration::from_secs(6))
.build()
.map_err(NsError::Client)?;
Ok(Self {
http,
base_url: site.base_url().to_string(),
token: site.token.clone(),
})
}
pub async fn verify_treatment_write(site: &Site) -> anyhow::Result<()> {
let token = site
.write_token
.as_deref()
.filter(|token| !token.trim().is_empty())
.ok_or_else(|| anyhow::anyhow!("no treatment write token configured"))?;
let http = reqwest::Client::builder()
.user_agent(concat!("sugarrush/", env!("CARGO_PKG_VERSION")))
.timeout(std::time::Duration::from_secs(12))
.connect_timeout(std::time::Duration::from_secs(6))
.build()
.map_err(NsError::Client)?;
let mut url = reqwest::Url::parse(&format!(
"{}/api/v2/authorization/request/",
site.base_url()
))?;
url.path_segments_mut()
.map_err(|_| anyhow::anyhow!("Nightscout URL cannot hold an authorization path"))?
.push(token);
let response = http
.get(url)
.send()
.await
.map_err(|e| NsError::Request("write-capability check failed", redact(e)))?;
let value: Value =
bounded_json(check_status(response)?, "invalid authorization response").await?;
if !has_treatment_create_permission(&value) {
anyhow::bail!("token is valid but lacks Nightscout treatment-create permission");
}
Ok(())
}
pub async fn create_treatment(
site: &Site,
body: &Value,
) -> std::result::Result<(), TreatmentWriteError> {
Self::verify_treatment_write(site)
.await
.map_err(|error| TreatmentWriteError::Definitive(error.to_string()))?;
let token = site.write_token.as_deref().unwrap_or_default();
let url = format!("{}/api/v1/treatments.json", site.base_url());
let response = reqwest::Client::builder()
.user_agent(concat!("sugarrush/", env!("CARGO_PKG_VERSION")))
.timeout(std::time::Duration::from_secs(12))
.connect_timeout(std::time::Duration::from_secs(6))
.build()
.map_err(|error| TreatmentWriteError::Definitive(error.to_string()))?
.post(url)
.query(&[("token", token)])
.json(body)
.send()
.await
.map_err(|_| {
TreatmentWriteError::Unknown(
"connection ended without a definitive Nightscout response".into(),
)
})?;
let status = response.status();
if !status.is_success() {
let message = format!("Nightscout returned HTTP {}", status.as_u16());
return Err(if status.is_client_error() {
TreatmentWriteError::Definitive(message)
} else {
TreatmentWriteError::Unknown(message)
});
}
let result: Value = bounded_json(response, "invalid treatment response")
.await
.map_err(|_| {
TreatmentWriteError::Unknown(
"Nightscout returned an unreadable success response".into(),
)
})?;
classify_v1_write(&result)
}
pub async fn entries_range(
&self,
start_ms: i64,
end_ms: i64,
want: usize,
) -> Result<Vec<Entry>> {
let count = want.max(1);
let url = format!("{}/api/v1/entries/sgv.json", self.base_url);
let resp = self
.http
.get(&url)
.query(&[
("find[date][$gte]", start_ms.to_string()),
("find[date][$lte]", end_ms.to_string()),
("count", count.to_string()),
("token", self.token.clone()),
])
.send()
.await
.map_err(|e| {
NsError::Request(
"can't reach Nightscout (check the URL and your connection)",
redact(e),
)
})?;
let entries: Vec<Entry> =
bounded_json(check_status(resp)?, "failed to parse Nightscout response").await?;
Ok(clean_newest_first(entries, end_ms))
}
pub async fn device_status(&self) -> Result<(DeviceStatus, Option<Vec<Prediction>>)> {
let url = format!("{}/api/v1/devicestatus.json", self.base_url);
let resp = self
.http
.get(&url)
.query(&[("count", "1"), ("token", self.token.as_str())])
.send()
.await
.map_err(|e| NsError::Request("devicestatus request failed", redact(e)))?;
let value: Value =
bounded_json(check_status(resp)?, "failed to parse devicestatus response").await?;
let latest = value.as_array().and_then(|items| items.first());
Ok((
latest.map(parse_device_status).unwrap_or_default(),
latest.and_then(parse_predicted),
))
}
pub async fn sensor_start(&self) -> Result<Option<i64>> {
let url = format!("{}/api/v1/treatments.json", self.base_url);
let resp = self
.http
.get(&url)
.query(&[
("find[eventType][$regex]", "Sensor"),
("count", "20"),
("token", self.token.as_str()),
])
.send()
.await
.map_err(|e| NsError::Request("treatments request failed", redact(e)))?;
let value: Value =
bounded_json(check_status(resp)?, "failed to parse treatments response").await?;
Ok(parse_sensor_start(&value))
}
pub async fn treatments(&self, start_ms: i64, end_ms: i64) -> Result<Vec<Treatment>> {
let since = DateTime::from_timestamp_millis(start_ms)
.map(|dt| dt.to_rfc3339())
.unwrap_or_default();
let url = format!("{}/api/v1/treatments.json", self.base_url);
let resp = self
.http
.get(&url)
.query(&[
("find[created_at][$gte]", since.as_str()),
("count", "300"),
("token", self.token.as_str()),
])
.send()
.await
.map_err(|e| NsError::Request("treatments request failed", redact(e)))?;
let value: Value =
bounded_json(check_status(resp)?, "failed to parse treatments response").await?;
Ok(parse_treatments(&value, start_ms, end_ms))
}
}
fn classify_v1_write(result: &Value) -> std::result::Result<(), TreatmentWriteError> {
match result.as_array() {
Some(created) if !created.is_empty() => Ok(()),
Some(_) => Err(TreatmentWriteError::Unknown(
"Nightscout accepted the request but reported no created treatment".into(),
)),
None => Err(TreatmentWriteError::Unknown(
"Nightscout returned an unexpected success body".into(),
)),
}
}
#[derive(Debug)]
pub enum TreatmentWriteError {
Definitive(String),
Unknown(String),
}
impl std::fmt::Display for TreatmentWriteError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
Self::Definitive(message) => write!(f, "treatment rejected: {message}"),
Self::Unknown(message) => write!(f, "treatment outcome unknown: {message}"),
}
}
}
impl std::error::Error for TreatmentWriteError {}
fn has_treatment_create_permission(value: &Value) -> bool {
value
.get("permissionGroups")
.and_then(Value::as_array)
.into_iter()
.flatten()
.flat_map(|group| group.as_array().into_iter().flatten())
.filter_map(Value::as_str)
.any(|permission| permission == "*" || permission == "api:treatments:create")
}
fn parse_sensor_start(value: &Value) -> Option<i64> {
value.as_array().and_then(|items| {
items
.iter()
.filter_map(|t| {
let event = t.get("eventType")?.as_str()?;
event
.contains("Sensor")
.then(|| {
t.get("created_at")
.and_then(Value::as_str)
.and_then(parse_iso)
})
.flatten()
})
.max()
})
}
fn parse_treatments(value: &Value, start_ms: i64, end_ms: i64) -> Vec<Treatment> {
value
.as_array()
.map(|items| {
items
.iter()
.filter_map(|t| {
let at = t.get("mills").and_then(Value::as_i64).or_else(|| {
t.get("created_at")
.and_then(Value::as_str)
.and_then(parse_iso)
})?;
if at < start_ms || at > end_ms {
return None;
}
let carbs = t.get("carbs").and_then(Value::as_f64).filter(|c| *c > 0.0);
let insulin = t
.get("insulin")
.and_then(Value::as_f64)
.filter(|i| *i > 0.0);
(carbs.is_some() || insulin.is_some()).then_some(Treatment {
at_ms: at,
carbs,
insulin,
})
})
.collect()
})
.unwrap_or_default()
}
#[derive(Debug, Clone)]
pub struct Treatment {
pub at_ms: i64,
pub carbs: Option<f64>,
pub insulin: Option<f64>,
}
#[derive(Debug, Clone, Default)]
pub struct DeviceStatus {
pub battery: Option<i64>,
pub device: Option<String>,
pub last_ms: Option<i64>,
pub iob: Option<f64>,
pub cob: Option<f64>,
}
fn parse_device_status(item: &Value) -> DeviceStatus {
let battery = item
.get("uploader")
.and_then(|u| u.get("battery"))
.or_else(|| item.get("uploaderBattery"))
.and_then(Value::as_i64);
let device = item
.get("device")
.and_then(Value::as_str)
.map(str::to_string);
let last_ms = item.get("mills").and_then(Value::as_i64).or_else(|| {
item.get("created_at")
.and_then(Value::as_str)
.and_then(parse_iso)
});
let l = item.get("loop");
let s = item.get("openaps").and_then(|o| o.get("suggested"));
let iob = l
.and_then(|l| l.get("iob"))
.and_then(|i| i.get("iob"))
.or_else(|| s.and_then(|s| s.get("IOB")))
.and_then(Value::as_f64);
let cob = l
.and_then(|l| l.get("cob"))
.and_then(|c| c.get("cob"))
.or_else(|| s.and_then(|s| s.get("COB")))
.and_then(Value::as_f64);
DeviceStatus {
battery,
device,
last_ms,
iob,
cob,
}
}
#[derive(Debug, Clone, Copy)]
pub struct Prediction {
pub at_ms: i64,
pub low: f64,
pub high: f64,
}
fn parse_predicted(item: &Value) -> Option<Vec<Prediction>> {
if let Some(pred) = item.get("loop").and_then(|l| l.get("predicted")) {
let start = pred
.get("startDate")
.and_then(Value::as_str)
.and_then(parse_iso)?;
let values = pred.get("values")?.as_array()?;
return Some(envelope(start, &[values]));
}
if let Some(sug) = item.get("openaps").and_then(|o| o.get("suggested")) {
let start = sug
.get("timestamp")
.and_then(Value::as_str)
.and_then(parse_iso)?;
let curves: Vec<&Vec<Value>> = sug
.get("predBGs")?
.as_object()?
.values()
.filter_map(Value::as_array)
.collect();
if curves.is_empty() {
return None;
}
return Some(envelope(start, &curves));
}
None
}
fn envelope(start_ms: i64, curves: &[&Vec<Value>]) -> Vec<Prediction> {
let max_len = curves.iter().map(|c| c.len()).max().unwrap_or(0);
let mut out = Vec::with_capacity(max_len);
for i in 0..max_len {
let mut lo = f64::MAX;
let mut hi = f64::MIN;
for c in curves {
if let Some(v) = c.get(i).and_then(Value::as_f64) {
lo = lo.min(v);
hi = hi.max(v);
}
}
if lo <= hi {
out.push(Prediction {
at_ms: start_ms + i as i64 * PRED_STEP_MS,
low: lo,
high: hi,
});
}
}
out
}
fn redact(e: reqwest::Error) -> reqwest::Error {
e.without_url()
}
fn check_status(resp: reqwest::Response) -> Result<reqwest::Response> {
use reqwest::StatusCode;
match resp.status() {
s if s.is_success() => Ok(resp),
StatusCode::UNAUTHORIZED | StatusCode::FORBIDDEN => Err(NsError::Auth),
StatusCode::NOT_FOUND => Err(NsError::NotFound),
s => Err(NsError::Http(s.as_u16())),
}
}
async fn bounded_json<T: DeserializeOwned>(
resp: reqwest::Response,
context: &'static str,
) -> Result<T> {
if resp
.content_length()
.is_some_and(|n| n > MAX_RESPONSE_BYTES)
{
return Err(NsError::TooLarge);
}
let bytes = resp
.bytes()
.await
.map_err(|e| NsError::Parse(context, redact(e)))?;
if bytes.len() as u64 > MAX_RESPONSE_BYTES {
return Err(NsError::TooLarge);
}
serde_json::from_slice(&bytes).map_err(|e| NsError::InvalidJson(context, e))
}
#[derive(Debug)]
pub enum NsError {
Client(reqwest::Error),
Request(&'static str, reqwest::Error),
Parse(&'static str, reqwest::Error),
InvalidJson(&'static str, serde_json::Error),
TooLarge,
Auth,
NotFound,
Http(u16),
}
impl NsError {
pub fn is_permanent(&self) -> bool {
matches!(self, Self::Auth | Self::NotFound)
}
}
impl std::fmt::Display for NsError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
NsError::Client(_) => write!(f, "failed to build HTTP client"),
NsError::Request(context, _)
| NsError::Parse(context, _)
| NsError::InvalidJson(context, _) => f.write_str(context),
NsError::TooLarge => write!(f, "Nightscout response exceeded the 8 MiB safety limit"),
NsError::Auth => write!(
f,
"authentication failed — check your read-only token (a Nightscout \
Subject token with the 'readable' role, not API_SECRET)"
),
NsError::NotFound => write!(
f,
"no Nightscout API at this URL (HTTP 404) — check the site URL"
),
NsError::Http(status) => write!(f, "Nightscout returned HTTP {status}"),
}
}
}
impl std::error::Error for NsError {
fn source(&self) -> Option<&(dyn std::error::Error + 'static)> {
match self {
Self::Client(e) | Self::Request(_, e) | Self::Parse(_, e) => Some(e),
Self::InvalidJson(_, e) => Some(e),
Self::Auth | Self::NotFound | Self::Http(_) | Self::TooLarge => None,
}
}
}
fn clean_newest_first(entries: Vec<Entry>, requested_end_ms: i64) -> Vec<Entry> {
let mut out: Vec<Entry> = entries
.into_iter()
.filter(|e| {
e.sgv >= MIN_PHYSIOLOGICAL_SGV
&& e.date <= requested_end_ms.saturating_add(MAX_FUTURE_SKEW_MS)
})
.collect();
out.sort_by_key(|e| std::cmp::Reverse(e.date));
out.dedup_by_key(|e| e.date);
out
}
fn parse_iso(s: &str) -> Option<i64> {
DateTime::parse_from_rfc3339(s)
.ok()
.map(|dt| dt.timestamp_millis())
}
#[cfg(test)]
pub mod fake {
use std::sync::Arc;
use tokio::io::{AsyncReadExt, AsyncWriteExt};
use tokio::net::TcpListener;
use crate::config::Site;
pub async fn serve(status: u16, body: &'static str) -> Site {
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
let port = listener.local_addr().unwrap().port();
let reason = match status {
200 => "OK",
401 => "Unauthorized",
403 => "Forbidden",
404 => "Not Found",
_ => "Internal Server Error",
};
let response = Arc::new(format!(
"HTTP/1.1 {status} {reason}\r\nContent-Type: application/json\r\n\
Content-Length: {}\r\nConnection: close\r\n\r\n{body}",
body.len()
));
tokio::spawn(async move {
loop {
let Ok((mut sock, _)) = listener.accept().await else {
return;
};
let response = Arc::clone(&response);
tokio::spawn(async move {
let mut buf = [0u8; 2048];
let _ = sock.read(&mut buf).await;
let _ = sock.write_all(response.as_bytes()).await;
let _ = sock.flush().await;
});
}
});
Site {
id: String::new(),
name: "fake".into(),
url: format!("http://127.0.0.1:{port}"),
token: "test-token".into(),
write_token: None,
timezone: None,
alerts: None,
}
}
pub async fn serve_slow(delay_ms: u64) -> Site {
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
let port = listener.local_addr().unwrap().port();
let body = "[]";
let response = Arc::new(format!(
"HTTP/1.1 200 OK\r\nContent-Type: application/json\r\n\
Content-Length: {}\r\nConnection: close\r\n\r\n{body}",
body.len()
));
tokio::spawn(async move {
loop {
let Ok((mut sock, _)) = listener.accept().await else {
return;
};
let response = Arc::clone(&response);
tokio::spawn(async move {
let mut buf = [0u8; 2048];
let _ = sock.read(&mut buf).await;
tokio::time::sleep(std::time::Duration::from_millis(delay_ms)).await;
let _ = sock.write_all(response.as_bytes()).await;
let _ = sock.flush().await;
});
}
});
Site {
id: String::new(),
name: "slow".into(),
url: format!("http://127.0.0.1:{port}"),
token: String::new(),
write_token: None,
timezone: None,
alerts: None,
}
}
pub async fn serve_stalled() -> Site {
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
let port = listener.local_addr().unwrap().port();
tokio::spawn(async move {
let mut held = Vec::new();
loop {
let Ok((sock, _)) = listener.accept().await else {
return;
};
held.push(sock); }
});
Site {
id: String::new(),
name: "stalled".into(),
url: format!("http://127.0.0.1:{port}"),
token: String::new(),
write_token: None,
timezone: None,
alerts: None,
}
}
}
#[cfg(test)]
mod tests {
use super::*;
fn entry(sgv: f64, date: i64) -> Entry {
Entry {
sgv,
date,
direction: None,
}
}
#[test]
fn a_v1_success_body_is_accepted_not_rejected() {
let created = serde_json::json!([{
"_id": "66b2f0a1c3d4e5f600112233",
"eventType": "Carb Correction",
"carbs": 15.0,
"enteredBy": "sugarrush"
}]);
assert!(
classify_v1_write(&created).is_ok(),
"the array Nightscout returns on success must not read as a rejection"
);
}
#[test]
fn an_unrecognised_success_body_is_unknown_not_a_verdict() {
for body in [
serde_json::json!([]),
serde_json::json!({"status": 200}),
serde_json::json!("ok"),
serde_json::Value::Null,
] {
match classify_v1_write(&body) {
Err(TreatmentWriteError::Unknown(_)) => {}
other => panic!("{body:?} should be Unknown, got {other:?}"),
}
}
}
#[test]
fn only_treatment_create_or_admin_crosses_the_write_boundary() {
assert!(has_treatment_create_permission(&serde_json::json!({
"permissionGroups": [["api:treatments:create"]]
})));
assert!(has_treatment_create_permission(&serde_json::json!({
"permissionGroups": [["*"]]
})));
assert!(!has_treatment_create_permission(&serde_json::json!({
"permissionGroups": [["*:*:read"]]
})));
}
#[test]
fn newest_first_regardless_of_server_order() {
let out = clean_newest_first(
vec![
entry(100.0, 1_000),
entry(120.0, 3_000),
entry(110.0, 2_000),
],
3_000,
);
assert_eq!(
out.iter().map(|e| e.date).collect::<Vec<_>>(),
vec![3_000, 2_000, 1_000]
);
assert_eq!(out[0].sgv, 120.0);
}
#[tokio::test]
async fn errors_never_carry_the_token() {
let site = Site {
id: String::new(),
name: "default".into(),
url: "http://127.0.0.1:1".into(),
token: "SEKRIT-TOKEN-9Q7X".into(),
write_token: None,
timezone: None,
alerts: None,
};
let client = Client::for_site(&site).unwrap();
let err = client.entries_range(0, 1, 1).await.unwrap_err();
let chain = format!("{err:#}\n{err:?}");
assert!(
!chain.contains("SEKRIT-TOKEN-9Q7X"),
"token leaked into the error chain:\n{chain}"
);
assert!(
!chain.contains("token="),
"query string leaked into the error chain:\n{chain}"
);
}
#[tokio::test]
async fn a_real_response_round_trips_over_the_wire() {
let body = r#"[{"sgv":142,"date":1700000000000,"direction":"Flat"},
{"sgv":138,"date":1699999700000,"direction":"FortyFiveDown"}]"#;
let site = fake::serve(200, body).await;
let client = Client::for_site(&site).unwrap();
let entries = client
.entries_range(0, 1_800_000_000_000, 10)
.await
.unwrap();
assert_eq!(entries.len(), 2);
assert_eq!(entries[0].sgv, 142.0); }
#[tokio::test]
async fn a_401_is_permanent_so_the_client_stops_retrying() {
let site = fake::serve(401, "unauthorized").await;
let client = Client::for_site(&site).unwrap();
let err = client.entries_range(0, 1, 1).await.unwrap_err();
assert!(matches!(&err, NsError::Auth));
assert!(err.is_permanent());
}
#[tokio::test]
async fn a_404_says_this_is_not_a_nightscout() {
let site = fake::serve(404, "nope").await;
let client = Client::for_site(&site).unwrap();
let err = client.entries_range(0, 1, 1).await.unwrap_err();
assert!(matches!(&err, NsError::NotFound));
assert!(err.is_permanent());
}
#[tokio::test]
async fn a_500_is_transient_so_the_backoff_keeps_trying() {
let site = fake::serve(500, "boom").await;
let client = Client::for_site(&site).unwrap();
let err = client.entries_range(0, 1, 1).await.unwrap_err();
assert!(matches!(&err, NsError::Http(500)));
assert!(!err.is_permanent());
assert!(err.to_string().contains("500"), "{err}");
}
#[tokio::test]
async fn a_proxy_serving_html_is_reported_as_a_parse_failure() {
let site = fake::serve(200, "<!doctype html><title>Sign in</title>").await;
let client = Client::for_site(&site).unwrap();
let err = client.entries_range(0, 1, 1).await.unwrap_err();
assert!(err.to_string().contains("parse"), "{err}");
}
#[tokio::test]
async fn an_empty_site_is_not_an_error() {
let site = fake::serve(200, "[]").await;
let client = Client::for_site(&site).unwrap();
assert!(client.entries_range(0, 1, 1).await.unwrap().is_empty());
}
#[tokio::test]
async fn oversized_responses_are_rejected_before_json_decode() {
let body: &'static str =
Box::leak("x".repeat(MAX_RESPONSE_BYTES as usize + 1).into_boxed_str());
let site = fake::serve(200, body).await;
let err = Client::for_site(&site)
.unwrap()
.entries_range(0, 1, 1)
.await
.unwrap_err();
assert!(matches!(err, NsError::TooLarge));
}
#[test]
fn entries_deserialize_from_a_nightscout_payload() {
let raw = r#"[
{"_id":"a","sgv":142,"date":1700000000000,"dateString":"2023-11-14T22:13:20.000Z",
"trend":4,"direction":"Flat","device":"share2","type":"sgv"},
{"_id":"b","sgv":138,"date":1699999700000,"direction":"FortyFiveDown","type":"sgv"}
]"#;
let entries: Vec<Entry> = serde_json::from_str(raw).unwrap();
assert_eq!(entries.len(), 2);
assert_eq!(entries[0].sgv, 142.0);
assert_eq!(entries[0].date, 1_700_000_000_000);
assert_eq!(entries[0].arrow(), "→");
assert_eq!(entries[1].arrow(), "↘");
}
#[test]
fn entry_survives_a_missing_direction() {
let e: Entry = serde_json::from_str(r#"{"sgv":95,"date":1}"#).unwrap();
assert!(e.direction.is_none());
assert_eq!(e.arrow(), "-");
}
#[test]
fn parses_a_loop_forecast() {
let raw = r#"{
"device":"loop://iPhone","created_at":"2023-11-14T22:10:00.000Z",
"loop":{"iob":{"iob":1.35},"cob":{"cob":18},
"predicted":{"startDate":"2023-11-14T22:10:00.000Z",
"values":[120,118,115,110]}}
}"#;
let item: Value = serde_json::from_str(raw).unwrap();
let preds = parse_predicted(&item).unwrap();
assert_eq!(preds.len(), 4);
assert_eq!(preds[0].low, 120.0);
assert_eq!(preds[0].high, 120.0);
assert_eq!(preds[3].low, 110.0);
assert_eq!(preds[1].at_ms - preds[0].at_ms, PRED_STEP_MS);
let status = parse_device_status(&item);
assert_eq!(status.iob, Some(1.35));
assert_eq!(status.cob, Some(18.0));
assert_eq!(status.device.as_deref(), Some("loop://iPhone"));
}
#[test]
fn parses_an_openaps_forecast_as_an_envelope() {
let raw = r#"{
"device":"openaps://rig","mills":1700000000000,
"openaps":{"suggested":{"timestamp":"2023-11-14T22:10:00.000Z","IOB":0.8,"COB":24,
"predBGs":{"IOB":[120,118,116],"ZT":[120,125,130],"UAM":[120,110,100]}}},
"uploader":{"battery":76}
}"#;
let item: Value = serde_json::from_str(raw).unwrap();
let preds = parse_predicted(&item).unwrap();
assert_eq!(preds.len(), 3);
assert_eq!(preds[0].low, 120.0);
assert_eq!(preds[0].high, 120.0);
assert_eq!(preds[2].low, 100.0);
assert_eq!(preds[2].high, 130.0);
let status = parse_device_status(&item);
assert_eq!(status.iob, Some(0.8));
assert_eq!(status.cob, Some(24.0));
assert_eq!(status.battery, Some(76));
assert_eq!(status.last_ms, Some(1_700_000_000_000));
}
#[test]
fn device_status_falls_back_to_uploader_battery_field() {
let item: Value = serde_json::from_str(
r#"{"uploaderBattery":42,"created_at":"2023-11-14T22:10:00.000Z"}"#,
)
.unwrap();
let status = parse_device_status(&item);
assert_eq!(status.battery, Some(42));
assert!(status.last_ms.is_some()); assert!(status.iob.is_none());
}
#[test]
fn no_forecast_when_the_uploader_publishes_none() {
let item: Value =
serde_json::from_str(r#"{"device":"share2","uploader":{"battery":90}}"#).unwrap();
assert!(parse_predicted(&item).is_none());
}
#[test]
fn treatments_are_filtered_to_the_window_and_to_real_doses() {
let raw = r#"[
{"eventType":"Meal Bolus","mills":1500,"carbs":45,"insulin":4.2},
{"eventType":"Correction Bolus","created_at":"1970-01-01T00:00:02.000Z","insulin":1.1},
{"eventType":"Note","mills":1600,"notes":"felt low"},
{"eventType":"Temp Basal","mills":1700,"carbs":0,"insulin":0},
{"eventType":"Meal Bolus","mills":99999,"carbs":30}
]"#;
let value: Value = serde_json::from_str(raw).unwrap();
let t = parse_treatments(&value, 1000, 5000);
assert_eq!(t.len(), 2);
assert_eq!(t[0].carbs, Some(45.0));
assert_eq!(t[0].insulin, Some(4.2));
assert_eq!(t[1].at_ms, 2000); assert_eq!(t[1].carbs, None);
}
#[test]
fn sensor_start_takes_the_newest_sensor_event() {
let raw = r#"[
{"eventType":"Sensor Start","created_at":"2023-11-01T08:00:00.000Z"},
{"eventType":"Site Change","created_at":"2023-11-12T08:00:00.000Z"},
{"eventType":"Sensor Change","created_at":"2023-11-11T07:30:00.000Z"}
]"#;
let value: Value = serde_json::from_str(raw).unwrap();
let start = parse_sensor_start(&value).unwrap();
assert_eq!(start, parse_iso("2023-11-11T07:30:00.000Z").unwrap());
let none: Value = serde_json::from_str(
r#"[{"eventType":"Site Change","created_at":"2023-11-12T08:00:00.000Z"}]"#,
)
.unwrap();
assert!(parse_sensor_start(&none).is_none());
}
#[test]
fn duplicate_timestamps_are_collapsed() {
let out = clean_newest_first(
vec![entry(120.0, 2_000), entry(120.0, 2_000), entry(95.0, 1_000)],
2_000,
);
assert_eq!(out.len(), 2);
assert_eq!(out[0].sgv - out[1].sgv, 25.0);
}
#[test]
fn drops_sensor_error_codes() {
let out = clean_newest_first(vec![entry(5.0, 2_000), entry(90.0, 1_000)], 2_000);
assert_eq!(out.len(), 1);
assert_eq!(out[0].sgv, 90.0);
}
#[test]
fn far_future_readings_cannot_suppress_staleness() {
let now = 1_000_000;
let out = clean_newest_first(vec![entry(100.0, now), entry(110.0, now + 86_400_000)], now);
assert_eq!(out.len(), 1);
assert_eq!(out[0].date, now);
}
}