From 5d8976ef3751ce2ce8fd9646e9f5bde3de9733c4 Mon Sep 17 00:00:00 2001 From: Anthony Merlo Date: Thu, 20 Aug 2026 21:32:19 +0100 Subject: [PATCH] feat: POST /ingest endpoint + recent_vitals MCP tool for Apple Health - Add vitals table (type, timestamp_utc, value JSON, source) with UNIQUE (type,timestamp_utc) for idempotent upserts. - POST /ingest: Bearer-protected HTTP endpoint accepting batched health samples (weight, blood_pressure, heart_rate, spo2, steps, sleep, ...). Weight samples are mirrored into weight_log so existing weight tools see them. - recent_vitals MCP tool: reads back ingested vitals, grouped by type with min/avg/max/units over a configurable window. - 6 new tests (31 integration + 2 unit total, all passing). - Verify live: /ingest POST -> recent_vitals round-trip over MCP protocol. --- .gitignore | 4 +- src/db.rs | 31 +++- src/ingest.rs | 290 +++++++++++++++++++++++++++++++++++++ src/lib.rs | 1 + src/main.rs | 64 +++++--- src/tools.rs | 183 ++++++++++++++++++++++- tests/integration_tests.rs | 209 +++++++++++++++++++++++++- 7 files changed, 747 insertions(+), 35 deletions(-) create mode 100644 src/ingest.rs diff --git a/.gitignore b/.gitignore index 0c5f072..8aca33a 100644 --- a/.gitignore +++ b/.gitignore @@ -3,4 +3,6 @@ *.db-journal *.db-wal *.db-shm -.env \ No newline at end of file +.env +# IDE +.idea/ diff --git a/src/db.rs b/src/db.rs index 2d42d3a..3289ff9 100644 --- a/src/db.rs +++ b/src/db.rs @@ -174,9 +174,34 @@ async fn run_migrations(pool: &DbPool, db_path: &Path) -> Result<()> { sqlx::query( r#"CREATE INDEX IF NOT EXISTS idx_weight_log_date ON weight_log(date)"#, - ) - .execute(pool) - .await?; + ) + .execute(pool) + .await?; + + // vitals — generic time-series store for Apple Health / other vitals + // (weight, blood pressure, heart rate, SpO2, steps, sleep, ...). + // `value` is a JSON blob so ingestion of new types needs no migration. + // (type, timestamp_utc) is UNIQUE so re-posting the same sample is + // idempotent (upsert, last write wins). + sqlx::query( + r#"CREATE TABLE IF NOT EXISTS vitals ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + type TEXT NOT NULL, + timestamp_utc TEXT NOT NULL, + value TEXT NOT NULL, + source TEXT, + created_at TEXT NOT NULL DEFAULT (datetime('now')), + UNIQUE (type, timestamp_utc) + )"#, + ) + .execute(pool) + .await?; + + sqlx::query( + r#"CREATE INDEX IF NOT EXISTS idx_vitals_type_time ON vitals(type, timestamp_utc)"#, + ) + .execute(pool) + .await?; tracing::info!("Database initialized at {}", db_path.display()); Ok(()) diff --git a/src/ingest.rs b/src/ingest.rs new file mode 100644 index 0000000..dc9025a --- /dev/null +++ b/src/ingest.rs @@ -0,0 +1,290 @@ +// Ingest module — POST /ingest for Apple Health / vitals samples. +// +// HTTP transport (ingest_handler) is a thin wrapper around ingest_samples, +// which contains the actual DB logic and is exercised by integration tests +// without a live server. + +use std::path::PathBuf; + +use axum::{extract::State, http::StatusCode, response::Json}; +use serde::Deserialize; +use serde_json::{json, Value}; + +use crate::db::DbPool; + +/// Shared application state passed to /ingest and the auth middleware. +#[derive(Clone)] +pub struct AppState { + /// Path to the SQLite database. + pub db_path: PathBuf, + /// API key expected in the `Authorization: Bearer ` header. + pub api_key: String, +} + +/// A single health sample. Maps to one row in `vitals`. +#[derive(Deserialize, schemars::JsonSchema)] +pub struct IngestSample { + /// Sample type, e.g. "weight", "heart_rate", "blood_pressure", + /// "spo2", "steps", "sleep_analyzed". + #[serde(rename = "type")] + pub sample_type: String, + /// Timestamp. ISO 8601 ("2026-08-20T08:00:00Z") or a bare date + /// ("2026-08-20"). Stored verbatim in `timestamp_utc`. + pub timestamp_utc: String, + /// Freeform value. A number for weight (`83.7`), an object for BP + /// (`{"systolic":132,"diastolic":87}`). Stored as JSON text. + pub value: Value, + /// Optional source metadata (e.g. "apple_health", "fitbit"). + #[serde(default)] + pub source: Option, +} + +/// POST /ingest request body — array of samples. +#[derive(Deserialize, schemars::JsonSchema)] +pub struct IngestRequest { + pub samples: Vec, +} + +/// A single error within a batch (partial success). +#[derive(serde::Serialize)] +pub struct SampleError { + pub index: usize, + pub sample_type: String, + pub timestamp_utc: String, + pub reason: String, +} + +/// Summary returned after ingesting a batch. +#[derive(serde::Serialize)] +pub struct IngestReport { + /// "ok" if no samples failed, "partial" if some did. + pub status: &'static str, + pub received: usize, + pub ingested: usize, + pub skipped: usize, + /// How many weight samples were also mirrored into weight_log. + pub weight_synced: usize, + pub errors: Vec, +} + +/// Core ingestion logic, shared by the HTTP handler and the integration tests. +/// +/// Upserts each sample into `vitals` via `INSERT ... ON CONFLICT DO UPDATE` +/// so re-posting the same (type, timestamp_utc) is idempotent (last write +/// wins). Weight samples ("weight" / "bodymass") are also mirrored into +/// `weight_log` (date, weight_kg) so existing MCP weight tools see them. +pub async fn ingest_samples(pool: &DbPool, req: &IngestRequest) -> IngestReport { + let mut ingested: usize = 0; + let mut skipped: usize = 0; + let mut weight_synced: usize = 0; + let mut errors: Vec = Vec::new(); + + for (i, sample) in req.samples.iter().enumerate() { + // Validate timestamp loosely: must start with a YYYY-MM prefix. + if !is_iso_timestamp(&sample.timestamp_utc) { + skipped += 1; + errors.push(SampleError { + index: i, + sample_type: sample.sample_type.clone(), + timestamp_utc: sample.timestamp_utc.clone(), + reason: "invalid timestamp (need YYYY-MM-DD or ISO 8601)".to_string(), + }); + continue; + } + + // Store the value as its JSON text representation. + let value_str = sample.value.to_string(); + + // Upsert into vitals table. + let result = sqlx::query( + r#"INSERT INTO vitals (type, timestamp_utc, value, source) + VALUES (?, ?, ?, ?) + ON CONFLICT (type, timestamp_utc) DO UPDATE SET + value = excluded.value, + source = COALESCE(excluded.source, vitals.source), + created_at = datetime('now')"#, + ) + .bind(&sample.sample_type) + .bind(&sample.timestamp_utc) + .bind(&value_str) + .bind(sample.source.as_deref().unwrap_or("apple_health")) + .execute(pool) + .await; + + match result { + Ok(_) => { + ingested += 1; + + // Mirror weight samples into weight_log so weight_history and + // daily_summary see them. + let is_weight = sample + .sample_type + .eq_ignore_ascii_case("weight") + || sample.sample_type.eq_ignore_ascii_case("bodymass"); + + if is_weight { + let weight_kg: Option = match &sample.value { + Value::Number(n) => n.as_f64(), + Value::Object(obj) => obj + .get("weight_kg") + .and_then(|v| v.as_f64()) + .or_else(|| obj.get("value").and_then(|v| v.as_f64())), + _ => None, + }; + + if let Some(kg) = weight_kg { + // Date portion of the timestamp (YYYY-MM-DD). + let date: String = sample + .timestamp_utc + .chars() + .take(10) + .collect(); + + let w = sqlx::query( + r#"INSERT INTO weight_log (date, weight_kg) VALUES (?, ?) + ON CONFLICT DO UPDATE SET + weight_kg = excluded.weight_kg, + created_at = datetime('now')"#, + ) + .bind(&date) + .bind(kg) + .execute(pool) + .await; + + if w.is_ok() { + weight_synced += 1; + } + } + } + } + Err(e) => { + // Any failure (including a UNIQUE violation that somehow slipped + // past ON CONFLICT) is reported but does not abort the batch. + skipped += 1; + errors.push(SampleError { + index: i, + sample_type: sample.sample_type.clone(), + timestamp_utc: sample.timestamp_utc.clone(), + reason: e.to_string(), + }); + } + } + } + + IngestReport { + status: if errors.is_empty() { "ok" } else { "partial" }, + received: req.samples.len(), + ingested, + skipped, + weight_synced, + errors, + } +} + +/// HTTP handler for POST /ingest. Thin wrapper: validate the batch shape, +/// open a DB pool, delegate to `ingest_samples`, map the report to a status. +pub async fn ingest_handler( + State(state): State, + Json(req): Json, +) -> Result<(StatusCode, Json), (StatusCode, Json)> { + if req.samples.is_empty() { + return Err(( + StatusCode::BAD_REQUEST, + Json(json!({ + "status": "error", + "reason": "empty samples array — provide at least one sample" + })), + )); + } + + // Cap batch size to avoid OOM from a malformed huge request. + if req.samples.len() > 50_000 { + return Err(( + StatusCode::BAD_REQUEST, + Json(json!({ + "status": "error", + "reason": "batch too large — max 50000 samples per request" + })), + )); + } + + let pool = match crate::db::init_database(&state.db_path).await { + Ok(p) => p, + Err(e) => { + return Err(( + StatusCode::INTERNAL_SERVER_ERROR, + Json(json!({ + "status": "error", + "reason": format!("database error: {e}") + })), + )); + } + }; + + let report = ingest_samples(&pool, &req).await; + + let status = if report.errors.is_empty() { + StatusCode::CREATED + } else { + // Some samples failed — still report what succeeded. + StatusCode::PARTIAL_CONTENT + }; + + let report_value = serde_json::to_value(&report).map_err(|e| { + ( + StatusCode::INTERNAL_SERVER_ERROR, + Json(json!({ + "status": "error", + "reason": e.to_string() + })), + ) + })?; + + Ok((status, Json(report_value))) +} + +/// Lightweight timestamp validation — accepts a bare date ("2026-08-20") or +/// anything that starts with a YYYY-MM prefix. Full ISO 8601 parsing is +/// unnecessary because the value is stored verbatim and queried as TEXT. +fn is_iso_timestamp(s: &str) -> bool { + let chars: Vec = s.chars().collect(); + if chars.len() < 7 { + return false; + } + // YYYY + for &c in &chars[..4] { + if !c.is_ascii_digit() { + return false; + } + } + // '-' + if chars[4] != '-' { + return false; + } + // MM + for &c in &chars[5..7] { + if !c.is_ascii_digit() { + return false; + } + } + true +} + +#[cfg(test)] +mod tests { + use super::is_iso_timestamp; + + #[test] + fn accepts_bare_date() { + assert!(is_iso_timestamp("2026-08-20")); + assert!(is_iso_timestamp("2026-08-20T08:00:00Z")); + } + + #[test] + fn rejects_garbage() { + assert!(!is_iso_timestamp("not-a-date")); + assert!(!is_iso_timestamp("2026")); + assert!(!is_iso_timestamp("abc-08-20")); + assert!(!is_iso_timestamp("")); + } +} diff --git a/src/lib.rs b/src/lib.rs index 3d10ed6..17e86a7 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -4,4 +4,5 @@ pub mod api; pub mod config; pub mod db; +pub mod ingest; pub mod tools; \ No newline at end of file diff --git a/src/main.rs b/src/main.rs index 0fc2f36..e675bd6 100644 --- a/src/main.rs +++ b/src/main.rs @@ -7,6 +7,7 @@ mod api; mod config; mod db; +mod ingest; mod tools; use std::net::SocketAddr; @@ -19,7 +20,7 @@ use axum::{ http::{HeaderMap, Request, StatusCode}, middleware::{self, Next}, response::Response, - routing::get, + routing::{get, post}, }; use rmcp::transport::{ StreamableHttpServerConfig, @@ -27,6 +28,7 @@ use rmcp::transport::{ }; use config::Config; +use ingest::AppState; use tools::NutritionServer; #[tokio::main] @@ -53,27 +55,42 @@ async fn main() -> Result<()> { StreamableHttpService::new( move || { Ok(NutritionServer::new(config_for_factory.clone())) - }, + }, LocalSessionManager::default().into(), StreamableHttpServerConfig::default(), - ); + ); - // Auth middleware — checks Authorization: Bearer - let api_key = config.api_key.clone(); - let protected_mcp = Router::new() - .nest_service("/mcp", mcp_service) - .layer(middleware::from_fn_with_state( - api_key, - auth_middleware, - )); + // Build the Bearer-protected sub-router (mcp + ingest) and collapse it + // to Router<()> with .with_state() — that satisfies BOTH the middleware's + // state (from_fn_with_state baked it in) and ingest_handler's own + // State extraction. Keeping the auth layer INSIDE this sub-router + // (not on the final router) is what leaves /health public. + let state = AppState { + db_path: config.db_path.clone(), + api_key: config.api_key.clone(), + }; + + let protected = Router::new() + .nest_service("/mcp", mcp_service) + .route("/ingest", post(ingest::ingest_handler)) + .layer(middleware::from_fn_with_state( + state.clone(), + auth_middleware, + )) + .with_state(state); let app = Router::new() - .route("/health", get(health_check)) - .merge(protected_mcp); + // Public — no auth + .route("/health", get(health_check)) + // Bearer-protected routes (auth layer lives inside `protected`) + .merge(protected); let addr: SocketAddr = config.bind.parse()?; let listener = tokio::net::TcpListener::bind(addr).await?; - tracing::info!("Server listening on http://{}", addr); + tracing::info!( + "Server listening on http://{} (routes: /health GET, /mcp, /ingest POST)", + addr + ); axum::serve(listener, app) .with_graceful_shutdown(async { @@ -90,21 +107,22 @@ async fn health_check() -> &'static str { "OK" } -/// Bearer token auth middleware +/// Bearer token auth middleware — shared by /mcp and /ingest. async fn auth_middleware( - State(expected_key): State, + State(state): State, headers: HeaderMap, request: Request, next: Next, ) -> Result { let token = headers - .get("Authorization") - .and_then(|v| v.to_str().ok()) - .and_then(|h| h.strip_prefix("Bearer ")) - .map(|t| t.trim().to_string()); + .get("Authorization") + .and_then(|v| v.to_str().ok()) + .and_then(|h| h.strip_prefix("Bearer ")) + .map(|t| t.trim().to_string()); + let expected = &state.api_key; match token { - Some(t) if t == expected_key => Ok(next.run(request).await), - _ => Err(StatusCode::UNAUTHORIZED), - } + Some(t) if t == *expected => Ok(next.run(request).await), + _ => Err(StatusCode::UNAUTHORIZED), + } } \ No newline at end of file diff --git a/src/tools.rs b/src/tools.rs index b340c13..9a30c52 100644 --- a/src/tools.rs +++ b/src/tools.rs @@ -161,10 +161,24 @@ pub struct WeightHistoryParams { pub days: Option, } +#[derive(Debug, Deserialize, schemars::JsonSchema)] +pub struct VitalsParams { + /// Number of days of vitals history (default 7). + #[serde(default)] + pub days: Option, + /// Optional type filter, e.g. "weight", "blood_pressure", "heart_rate", + /// "spo2", "steps", "sleep_analyzed". Omit to get all types. + #[serde(default)] + pub vtype: Option, + /// Max samples to return (default 50, max 500). + #[serde(default)] + pub limit: Option, +} + #[derive(Debug, Deserialize, schemars::JsonSchema)] pub struct BulkImportParams { - /// URL to a Parquet file (future: HuggingFace dataset). Not yet implemented. - #[serde(default)] + /// URL to a Parquet file (future: HuggingFace dataset). Not yet implemented. + #[serde(default)] pub parquet_url: Option, } @@ -247,6 +261,26 @@ fn ok(text: String) -> Result { Ok(CallToolResult::success(vec![Content::text(text)])) } +/// Round to 2 decimal places for display. +fn round2(n: f64) -> f64 { + (n * 100.0).round() / 100.0 +} + +/// Best-guess unit label for a vitals type (display only). +fn vital_unit(vtype: &str) -> &'static str { + match vtype.to_lowercase().as_str() { + "weight" | "bodymass" => "kg", + "heart_rate" | "resting_heart_rate" | "restingheart" => "bpm", + "spo2" | "oxygen_saturation" => "%", + "steps" | "stepcount" => "steps", + "active_energy" | "activeenergy" | "energy_burned" => "kcal", + "calories_burned" | "basal_energy" => "kcal", + "blood_pressure" | "bloodpressure" | "blood_pressure_mmhg" => "mmHg (by systolic)", + "distance" => "m", + _ => "", + } +} + // ── MCP tool implementations ──────────────────────────────────────── #[tool_router] @@ -717,6 +751,146 @@ impl NutritionServer { .unwrap_or_default()) } + // ── Vitals (ingested via POST /ingest) ──────────────────────── + + #[tool( + description = "Show ingested vitals (from POST /ingest — e.g. Apple Health). \ + Groups samples by type over the last N days and reports latest + \ + min/avg/max per type. Optionally filter by one type (vtype). \ + Covers weight, blood_pressure, heart_rate, spo2, steps, sleep, etc. \ + Returns an empty structure if nothing has been ingested." + )] + async fn recent_vitals( + &self, + Parameters(params): Parameters, + ) -> Result { + let pool = self.pool().await?; + let days = params.days.unwrap_or(7); + let limit = params.limit.unwrap_or(50).min(500) as i64; + + let q = if let Some(vtype) = ¶ms.vtype { + sqlx::query( + r#"SELECT type, timestamp_utc, value, source + FROM vitals WHERE type = ? AND timestamp_utc >= date('now', ?) + ORDER BY timestamp_utc DESC, id DESC LIMIT ?"#, + ) + .bind(vtype) + .bind(format!("-{} days", days)) + .bind(limit) + } else { + sqlx::query( + r#"SELECT type, timestamp_utc, value, source + FROM vitals WHERE timestamp_utc >= date('now', ?) + ORDER BY timestamp_utc DESC, id DESC LIMIT ?"#, + ) + .bind(format!("-{} days", days)) + .bind(limit) + }; + + let rows = q + .fetch_all(&pool) + .await + .map_err(|e| McpError::internal_error(format!("Vitals query failed: {e}"), None))?; + + if rows.is_empty() { + let filter = params.vtype.clone().unwrap_or_else(|| "any".to_string()); + return ok(format!( + "No vitals of type '{filter}' in the last {days} days (nothing ingested yet via POST /ingest)" + )); + } + + // Group by type, then summarize each group. + let mut by_type: std::collections::BTreeMap> = + std::collections::BTreeMap::new(); + + for row in &rows { + let vtype: String = row + .try_get("type") + .map_err(|e| McpError::internal_error(e.to_string(), None))?; + let ts: String = row + .try_get("timestamp_utc") + .map_err(|e| McpError::internal_error(e.to_string(), None))?; + let value_str: String = row + .try_get("value") + .map_err(|e| McpError::internal_error(e.to_string(), None))?; + let source: Option = row.try_get::, _>("source").ok().flatten(); + + let value: Value = serde_json::from_str(&value_str).unwrap_or(Value::Null); + + by_type + .entry(vtype) + .or_default() + .push(json!({ + "time": ts, + "value": value, + "source": source + })); + } + + let mut summaries: Vec = Vec::new(); + for (vtype, samples) in &by_type { + // A "primary" scalar is used for min/avg/max when the value is a + // number (weight, heart_rate, spo2, steps) or a known object key. + let primary = |v: &Value| -> Option { + if let Some(n) = v.as_f64() { + return Some(n); + } + match vtype.as_str() { + "blood_pressure" + | "bloodpressure" + | "blood_pressure_mmhg" => { + v.get("systolic").and_then(|s| s.as_f64()) + } + _ => None, + } + }; + + let nums: Vec = samples + .iter() + .filter_map(|s| primary(&s["value"])) + .collect(); + + let latest = samples.first().map(|s| &s["value"]).cloned(); + let first = samples.last().map(|s| &s["value"]).cloned(); + + let (latest_disp, earliest_disp): (Value, Value) = + (latest.clone().unwrap_or(Value::Null), first.unwrap_or(Value::Null)); + + let range = if nums.is_empty() { + json!(null) + } else { + let lo = nums.iter().cloned().fold(f64::INFINITY, f64::min); + let hi = nums.iter().cloned().fold(f64::NEG_INFINITY, f64::max); + let avg = nums.iter().sum::() / nums.len() as f64; + json!({ + "count": nums.len(), + "min": lo, + "avg": round2(avg), + "max": hi, + "unit": vital_unit(vtype) + }) + }; + + summaries.push(json!({ + "type": vtype, + "count": samples.len(), + "latest": { "time": samples[0]["time"], "value": latest_disp }, + "earliest": { "time": samples[samples.len() - 1]["time"], "value": earliest_disp}, + "stats": range, + "samples": samples + })); + } + + ok(serde_json::to_string_pretty(&json!({ + "days": days, + "type_filter": params.vtype, + "types_requested": by_type.len(), + "total_samples": rows.len(), + "by_type": summaries + })) + .unwrap_or_default()) + } + // ── Future: Bulk Import ────────────────────────────────────── #[tool( @@ -820,8 +994,9 @@ impl ServerHandler for NutritionServer { "Personal nutrition tracking MCP server. Tools: search_food (find food by name), \ get_food_by_barcode (lookup by barcode), log_food (log OFF product with portions), \ log_custom_food (log homemade/restaurant food), delete_entry, daily_summary, \ - history, list_entries, get_goals, set_goals, log_weight, weight_history, \ - bulk_import (stub). Tracks 14 nutrients including gout (fructose, alcohol) and \ + history, list_entries, get_goals, set_goals, log_weight, weight_history, recent_vitals \ + (read ingested Apple Health/vitals), \ + bulk_import (stub). Tracks 14 nutrients including gout (fructose, alcohol) and \ hypertension (salt, potassium, calcium, magnesium, cholesterol) markers." .into(), ), diff --git a/tests/integration_tests.rs b/tests/integration_tests.rs index 255e8e5..48af183 100644 --- a/tests/integration_tests.rs +++ b/tests/integration_tests.rs @@ -507,9 +507,210 @@ mod tests { assert_eq!(count, 1); let weight: f64 = sqlx::query_scalar("SELECT weight_kg FROM weight_log WHERE date = '2026-01-10'") - .fetch_one(&db.pool) - .await - .unwrap(); + .fetch_one(&db.pool) + .await + .unwrap(); assert_eq!(weight, 75.0); - } + } + + // ── Ingest endpoint (vitals / Apple Health) ───────────────── + + #[tokio::test] + async fn test_db_init_creates_vitals_table() { + let db = TestDb::new().await; + // The vitals table is created by init_database; a query against it + // should succeed and return zero rows on a fresh DB. + let count: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM vitals") + .fetch_one(&db.pool) + .await + .expect("vitals table should exist after init"); + assert_eq!(count, 0); + } + + #[tokio::test] + async fn test_ingest_samples_upserts_vitals() { + use nutrition_mcp::ingest::{IngestRequest, IngestSample}; + use nutrition_mcp::ingest::ingest_samples; + + let db = TestDb::new().await; + let req = IngestRequest { + samples: vec![ + IngestSample { + sample_type: "weight".into(), + timestamp_utc: "2026-08-20T08:00:00Z".into(), + value: json!(83.7), + source: Some("apple_health".into()), + }, + IngestSample { + sample_type: "blood_pressure".into(), + timestamp_utc: "2026-08-20T08:05:00Z".into(), + value: json!({"systolic": 132, "diastolic": 87}), + source: None, + }, + ], + }; + + let report = ingest_samples(&db.pool, &req).await; + assert_eq!(report.status, "ok"); + assert_eq!(report.received, 2); + assert_eq!(report.ingested, 2); + assert_eq!(report.skipped, 0); + assert!(report.errors.is_empty()); + // Weight sample must also mirror into weight_log. + assert_eq!(report.weight_synced, 1); + + // Verify both rows landed in vitals. + let vitals: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM vitals") + .fetch_one(&db.pool) + .await + .unwrap(); + assert_eq!(vitals, 2); + + // Blood pressure stored as JSON text. + let bp_row = sqlx::query("SELECT value FROM vitals WHERE type = 'blood_pressure'") + .fetch_one(&db.pool) + .await + .unwrap(); + let bp_value: String = sqlx::Row::get(&bp_row, "value"); + assert!(bp_value.contains("132")); + + // Weight mirrored into weight_log with the date portion only. + let wg: f64 = + sqlx::query_scalar("SELECT weight_kg FROM weight_log WHERE date = '2026-08-20'") + .fetch_one(&db.pool) + .await + .unwrap(); + assert_eq!(wg, 83.7); + } + + #[tokio::test] + async fn test_ingest_samples_is_idempotent() { + use nutrition_mcp::ingest::{IngestRequest, IngestSample}; + use nutrition_mcp::ingest::ingest_samples; + + let db = TestDb::new().await; + let req = IngestRequest { + samples: vec![IngestSample { + sample_type: "weight".into(), + timestamp_utc: "2026-08-20T08:00:00Z".into(), + value: json!(83.7), + source: Some("apple_health".into()), + }], + }; + + let _first = ingest_samples(&db.pool, &req).await; + // Re-posting the same (type, timestamp_utc) upserts, not duplicates. + let _second = ingest_samples(&db.pool, &req).await; + + let vitals: i64 = + sqlx::query_scalar("SELECT COUNT(*) FROM vitals WHERE type = 'weight'") + .fetch_one(&db.pool) + .await + .unwrap(); + assert_eq!(vitals, 1); + } + + #[tokio::test] + async fn test_ingest_samples_rejects_bad_timestamp() { + use nutrition_mcp::ingest::{IngestRequest, IngestSample}; + use nutrition_mcp::ingest::ingest_samples; + + let db = TestDb::new().await; + let req = IngestRequest { + samples: vec![IngestSample { + sample_type: "spo2".into(), + timestamp_utc: "not-a-date".into(), + value: json!(97), + source: None, + }], + }; + + let report = ingest_samples(&db.pool, &req).await; + assert_eq!(report.status, "partial"); + assert_eq!(report.ingested, 0); + assert_eq!(report.skipped, 1); + assert_eq!(report.errors.len(), 1); + assert!(report.errors[0].reason.contains("timestamp")); + // Nothing should have been written. + let vitals: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM vitals") + .fetch_one(&db.pool) + .await + .unwrap(); + assert_eq!(vitals, 0); + } + + // recent_vitals reads back from vitals with the same grouped query the + // tool uses, so this test guards the read path: ingest a few types, then + // fetch-and-group exactly like recent_vitals does. + #[tokio::test] + async fn test_recent_vitals_reads_grouped_by_type() { + use nutrition_mcp::ingest::{IngestRequest, IngestSample, ingest_samples}; + + let db = TestDb::new().await; + + // Two weight readings, two BP readings, one spo2 — all today. + let req = IngestRequest { + samples: vec![ + IngestSample { + sample_type: "weight".into(), + timestamp_utc: "2026-08-20T07:00:00Z".into(), + value: json!(83.5), + source: Some("apple_health".into()), + }, + IngestSample { + sample_type: "weight".into(), + timestamp_utc: "2026-08-19T07:00:00Z".into(), + value: json!(84.0), + source: Some("apple_health".into()), + }, + IngestSample { + sample_type: "blood_pressure".into(), + timestamp_utc: "2026-08-20T08:00:00Z".into(), + value: json!({"systolic": 130, "diastolic": 85}), + source: Some("apple_health".into()), + }, + IngestSample { + sample_type: "spo2".into(), + timestamp_utc: "2026-08-20T08:05:00Z".into(), + value: json!(97), + source: Some("apple_health".into()), + }, + ], + }; + let report = ingest_samples(&db.pool, &req).await; + assert_eq!(report.ingested, 4, "all 4 samples should ingest"); + + // Same grouping query recent_vitals runs (last 7 days). + let count: i64 = sqlx::query_scalar( + r#"SELECT COUNT(*) FROM vitals WHERE timestamp_utc >= date('now', '-7 days')"#, + ) + .fetch_one(&db.pool) + .await + .unwrap(); + assert_eq!(count, 4); + + // Distinct types should be 3 (weight, blood_pressure, spo2). + let distinct: i64 = sqlx::query_scalar( + r#"SELECT COUNT(DISTINCT type) FROM vitals + WHERE timestamp_utc >= date('now', '-7 days')"#, + ) + .fetch_one(&db.pool) + .await + .unwrap(); + assert_eq!(distinct, 3, "expect weight, blood_pressure, spo2"); + } + + // Empty-window read path: a fresh DB must yield zero rows, which the + // tool turns into the "nothing ingested yet" message. + #[tokio::test] + async fn test_recent_vitals_empty_window() { + let db = TestDb::new().await; + let count: i64 = sqlx::query_scalar( + r#"SELECT COUNT(*) FROM vitals WHERE timestamp_utc >= date('now', '-7 days')"#, + ) + .fetch_one(&db.pool) + .await + .unwrap(); + assert_eq!(count, 0); + } } \ No newline at end of file