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.
This commit is contained in:
Anthony Merlo 2026-08-20 21:32:19 +01:00
parent 171c8a0490
commit 5d8976ef37
7 changed files with 747 additions and 35 deletions

2
.gitignore vendored
View file

@ -4,3 +4,5 @@
*.db-wal *.db-wal
*.db-shm *.db-shm
.env .env
# IDE
.idea/

View file

@ -178,6 +178,31 @@ async fn run_migrations(pool: &DbPool, db_path: &Path) -> Result<()> {
.execute(pool) .execute(pool)
.await?; .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()); tracing::info!("Database initialized at {}", db_path.display());
Ok(()) Ok(())
} }

290
src/ingest.rs Normal file
View file

@ -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 <key>` 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<String>,
}
/// POST /ingest request body — array of samples.
#[derive(Deserialize, schemars::JsonSchema)]
pub struct IngestRequest {
pub samples: Vec<IngestSample>,
}
/// 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<SampleError>,
}
/// 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<SampleError> = 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<f64> = 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<AppState>,
Json(req): Json<IngestRequest>,
) -> Result<(StatusCode, Json<Value>), (StatusCode, Json<Value>)> {
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<char> = 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(""));
}
}

View file

@ -4,4 +4,5 @@
pub mod api; pub mod api;
pub mod config; pub mod config;
pub mod db; pub mod db;
pub mod ingest;
pub mod tools; pub mod tools;

View file

@ -7,6 +7,7 @@
mod api; mod api;
mod config; mod config;
mod db; mod db;
mod ingest;
mod tools; mod tools;
use std::net::SocketAddr; use std::net::SocketAddr;
@ -19,7 +20,7 @@ use axum::{
http::{HeaderMap, Request, StatusCode}, http::{HeaderMap, Request, StatusCode},
middleware::{self, Next}, middleware::{self, Next},
response::Response, response::Response,
routing::get, routing::{get, post},
}; };
use rmcp::transport::{ use rmcp::transport::{
StreamableHttpServerConfig, StreamableHttpServerConfig,
@ -27,6 +28,7 @@ use rmcp::transport::{
}; };
use config::Config; use config::Config;
use ingest::AppState;
use tools::NutritionServer; use tools::NutritionServer;
#[tokio::main] #[tokio::main]
@ -58,22 +60,37 @@ async fn main() -> Result<()> {
StreamableHttpServerConfig::default(), StreamableHttpServerConfig::default(),
); );
// Auth middleware — checks Authorization: Bearer <key> // Build the Bearer-protected sub-router (mcp + ingest) and collapse it
let api_key = config.api_key.clone(); // to Router<()> with .with_state() — that satisfies BOTH the middleware's
let protected_mcp = Router::new() // state (from_fn_with_state baked it in) and ingest_handler's own
// State<AppState> 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) .nest_service("/mcp", mcp_service)
.route("/ingest", post(ingest::ingest_handler))
.layer(middleware::from_fn_with_state( .layer(middleware::from_fn_with_state(
api_key, state.clone(),
auth_middleware, auth_middleware,
)); ))
.with_state(state);
let app = Router::new() let app = Router::new()
// Public — no auth
.route("/health", get(health_check)) .route("/health", get(health_check))
.merge(protected_mcp); // Bearer-protected routes (auth layer lives inside `protected`)
.merge(protected);
let addr: SocketAddr = config.bind.parse()?; let addr: SocketAddr = config.bind.parse()?;
let listener = tokio::net::TcpListener::bind(addr).await?; 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) axum::serve(listener, app)
.with_graceful_shutdown(async { .with_graceful_shutdown(async {
@ -90,9 +107,9 @@ async fn health_check() -> &'static str {
"OK" "OK"
} }
/// Bearer token auth middleware /// Bearer token auth middleware — shared by /mcp and /ingest.
async fn auth_middleware( async fn auth_middleware(
State(expected_key): State<String>, State(state): State<AppState>,
headers: HeaderMap, headers: HeaderMap,
request: Request<axum::body::Body>, request: Request<axum::body::Body>,
next: Next, next: Next,
@ -103,8 +120,9 @@ async fn auth_middleware(
.and_then(|h| h.strip_prefix("Bearer ")) .and_then(|h| h.strip_prefix("Bearer "))
.map(|t| t.trim().to_string()); .map(|t| t.trim().to_string());
let expected = &state.api_key;
match token { match token {
Some(t) if t == expected_key => Ok(next.run(request).await), Some(t) if t == *expected => Ok(next.run(request).await),
_ => Err(StatusCode::UNAUTHORIZED), _ => Err(StatusCode::UNAUTHORIZED),
} }
} }

View file

@ -161,6 +161,20 @@ pub struct WeightHistoryParams {
pub days: Option<u32>, pub days: Option<u32>,
} }
#[derive(Debug, Deserialize, schemars::JsonSchema)]
pub struct VitalsParams {
/// Number of days of vitals history (default 7).
#[serde(default)]
pub days: Option<u32>,
/// Optional type filter, e.g. "weight", "blood_pressure", "heart_rate",
/// "spo2", "steps", "sleep_analyzed". Omit to get all types.
#[serde(default)]
pub vtype: Option<String>,
/// Max samples to return (default 50, max 500).
#[serde(default)]
pub limit: Option<u32>,
}
#[derive(Debug, Deserialize, schemars::JsonSchema)] #[derive(Debug, Deserialize, schemars::JsonSchema)]
pub struct BulkImportParams { pub struct BulkImportParams {
/// URL to a Parquet file (future: HuggingFace dataset). Not yet implemented. /// URL to a Parquet file (future: HuggingFace dataset). Not yet implemented.
@ -247,6 +261,26 @@ fn ok(text: String) -> Result<CallToolResult, McpError> {
Ok(CallToolResult::success(vec![Content::text(text)])) 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 ──────────────────────────────────────── // ── MCP tool implementations ────────────────────────────────────────
#[tool_router] #[tool_router]
@ -717,6 +751,146 @@ impl NutritionServer {
.unwrap_or_default()) .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<VitalsParams>,
) -> Result<CallToolResult, McpError> {
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) = &params.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<String, Vec<Value>> =
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<String> = row.try_get::<Option<String>, _>("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<Value> = 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<f64> {
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<f64> = 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::<f64>() / 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 ────────────────────────────────────── // ── Future: Bulk Import ──────────────────────────────────────
#[tool( #[tool(
@ -820,7 +994,8 @@ impl ServerHandler for NutritionServer {
"Personal nutrition tracking MCP server. Tools: search_food (find food by name), \ "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), \ 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, \ log_custom_food (log homemade/restaurant food), delete_entry, daily_summary, \
history, list_entries, get_goals, set_goals, log_weight, weight_history, \ 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 \ bulk_import (stub). Tracks 14 nutrients including gout (fructose, alcohol) and \
hypertension (salt, potassium, calcium, magnesium, cholesterol) markers." hypertension (salt, potassium, calcium, magnesium, cholesterol) markers."
.into(), .into(),

View file

@ -512,4 +512,205 @@ mod tests {
.unwrap(); .unwrap();
assert_eq!(weight, 75.0); 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);
}
} }