498 lines
20 KiB
Python
498 lines
20 KiB
Python
"""
|
|
Aggregation Service - Computes aggregated metrics from raw Health Connect data.
|
|
|
|
This service calculates daily, weekly, and monthly aggregations on-demand
|
|
from the raw health records stored in the database.
|
|
"""
|
|
from sqlalchemy.orm import Session
|
|
from sqlalchemy import func, and_
|
|
from datetime import datetime, date, timedelta
|
|
from typing import Optional, List, Dict, Any
|
|
from dataclasses import dataclass
|
|
import logging
|
|
|
|
logger = logging.getLogger(__name__)
|
|
|
|
from backend.models.health_connect_raw import (
|
|
HCRawSteps,
|
|
HCRawHeartRate,
|
|
HCRawHeartRateSample,
|
|
HCRawRestingHeartRate,
|
|
HCRawSleepSession,
|
|
HCRawSleepStage,
|
|
HCRawDistance,
|
|
HCRawCalories,
|
|
HCRawOxygenSaturation,
|
|
HCRawWeight,
|
|
HCRawExerciseSession,
|
|
)
|
|
|
|
|
|
@dataclass
|
|
class DailyMetrics:
|
|
"""Aggregated daily health metrics"""
|
|
date: date
|
|
step_count: int = 0
|
|
distance_meters: float = 0.0
|
|
calories_burned: float = 0.0
|
|
active_calories: float = 0.0
|
|
|
|
# Sleep
|
|
sleep_duration_minutes: int = 0
|
|
deep_sleep_minutes: int = 0
|
|
light_sleep_minutes: int = 0
|
|
rem_sleep_minutes: int = 0
|
|
awake_minutes: int = 0
|
|
|
|
# Heart rate
|
|
avg_heart_rate: Optional[int] = None
|
|
min_heart_rate: Optional[int] = None
|
|
max_heart_rate: Optional[int] = None
|
|
resting_heart_rate: Optional[int] = None
|
|
|
|
# SpO2
|
|
avg_spo2: Optional[float] = None
|
|
min_spo2: Optional[float] = None
|
|
max_spo2: Optional[float] = None
|
|
|
|
# Body
|
|
weight: Optional[float] = None
|
|
|
|
# Counts
|
|
exercise_count: int = 0
|
|
|
|
def to_dict(self) -> Dict[str, Any]:
|
|
return {
|
|
"date": self.date.isoformat(),
|
|
"step_count": self.step_count,
|
|
"distance_meters": self.distance_meters,
|
|
"calories_burned": self.calories_burned,
|
|
"active_calories": self.active_calories,
|
|
"sleep_duration_minutes": self.sleep_duration_minutes,
|
|
"deep_sleep_minutes": self.deep_sleep_minutes,
|
|
"light_sleep_minutes": self.light_sleep_minutes,
|
|
"rem_sleep_minutes": self.rem_sleep_minutes,
|
|
"awake_minutes": self.awake_minutes,
|
|
"avg_heart_rate": self.avg_heart_rate,
|
|
"min_heart_rate": self.min_heart_rate,
|
|
"max_heart_rate": self.max_heart_rate,
|
|
"resting_heart_rate": self.resting_heart_rate,
|
|
"avg_spo2": self.avg_spo2,
|
|
"min_spo2": self.min_spo2,
|
|
"max_spo2": self.max_spo2,
|
|
"weight": self.weight,
|
|
"exercise_count": self.exercise_count,
|
|
}
|
|
|
|
|
|
class AggregationService:
|
|
"""
|
|
Service for computing aggregated metrics from raw Health Connect data.
|
|
|
|
All methods take a user_id and return computed metrics based on the
|
|
raw records stored in the database.
|
|
"""
|
|
|
|
def __init__(self, db: Session):
|
|
self.db = db
|
|
|
|
# ==================== DAILY AGGREGATIONS ====================
|
|
|
|
def get_daily_steps(self, user_id: int, target_date: date) -> int:
|
|
"""Sum all step records for a given day."""
|
|
start_of_day = datetime.combine(target_date, datetime.min.time())
|
|
end_of_day = datetime.combine(target_date + timedelta(days=1), datetime.min.time())
|
|
|
|
result = self.db.query(func.sum(HCRawSteps.count)).filter(
|
|
HCRawSteps.user_id == user_id,
|
|
HCRawSteps.start_time >= start_of_day,
|
|
HCRawSteps.start_time < end_of_day
|
|
).scalar()
|
|
|
|
return result or 0
|
|
|
|
def get_daily_distance(self, user_id: int, target_date: date) -> float:
|
|
"""Sum all distance records for a given day."""
|
|
start_of_day = datetime.combine(target_date, datetime.min.time())
|
|
end_of_day = datetime.combine(target_date + timedelta(days=1), datetime.min.time())
|
|
|
|
result = self.db.query(func.sum(HCRawDistance.distance_meters)).filter(
|
|
HCRawDistance.user_id == user_id,
|
|
HCRawDistance.start_time >= start_of_day,
|
|
HCRawDistance.start_time < end_of_day
|
|
).scalar()
|
|
|
|
return result or 0.0
|
|
|
|
def get_daily_calories(self, user_id: int, target_date: date) -> Dict[str, float]:
|
|
"""Get calories burned for a given day (active and total)."""
|
|
start_of_day = datetime.combine(target_date, datetime.min.time())
|
|
end_of_day = datetime.combine(target_date + timedelta(days=1), datetime.min.time())
|
|
|
|
# Active calories
|
|
active = self.db.query(func.sum(HCRawCalories.calories)).filter(
|
|
HCRawCalories.user_id == user_id,
|
|
HCRawCalories.start_time >= start_of_day,
|
|
HCRawCalories.start_time < end_of_day,
|
|
HCRawCalories.record_type == "active"
|
|
).scalar() or 0.0
|
|
|
|
# Total calories
|
|
total = self.db.query(func.sum(HCRawCalories.calories)).filter(
|
|
HCRawCalories.user_id == user_id,
|
|
HCRawCalories.start_time >= start_of_day,
|
|
HCRawCalories.start_time < end_of_day,
|
|
HCRawCalories.record_type == "total"
|
|
).scalar() or 0.0
|
|
|
|
# If no total, use active as total
|
|
if total == 0 and active > 0:
|
|
total = active
|
|
|
|
return {"active": active, "total": total}
|
|
|
|
def get_daily_heart_rate_stats(self, user_id: int, target_date: date) -> Dict[str, Optional[int]]:
|
|
"""Calculate min, max, avg heart rate for a day from all samples."""
|
|
start_of_day = datetime.combine(target_date, datetime.min.time())
|
|
end_of_day = datetime.combine(target_date + timedelta(days=1), datetime.min.time())
|
|
|
|
# Get all heart rate records for the day
|
|
hr_records = self.db.query(HCRawHeartRate).filter(
|
|
HCRawHeartRate.user_id == user_id,
|
|
HCRawHeartRate.start_time >= start_of_day,
|
|
HCRawHeartRate.start_time < end_of_day
|
|
).all()
|
|
|
|
if not hr_records:
|
|
return {"avg": None, "min": None, "max": None}
|
|
|
|
# Collect all samples
|
|
all_bpm = []
|
|
for record in hr_records:
|
|
for sample in record.samples:
|
|
if sample.bpm and sample.bpm > 0:
|
|
all_bpm.append(sample.bpm)
|
|
|
|
if not all_bpm:
|
|
return {"avg": None, "min": None, "max": None}
|
|
|
|
return {
|
|
"avg": int(sum(all_bpm) / len(all_bpm)),
|
|
"min": min(all_bpm),
|
|
"max": max(all_bpm)
|
|
}
|
|
|
|
def get_daily_resting_heart_rate(self, user_id: int, target_date: date) -> Optional[int]:
|
|
"""Get the resting heart rate for a day (last recorded value)."""
|
|
start_of_day = datetime.combine(target_date, datetime.min.time())
|
|
end_of_day = datetime.combine(target_date + timedelta(days=1), datetime.min.time())
|
|
|
|
record = self.db.query(HCRawRestingHeartRate).filter(
|
|
HCRawRestingHeartRate.user_id == user_id,
|
|
HCRawRestingHeartRate.time >= start_of_day,
|
|
HCRawRestingHeartRate.time < end_of_day
|
|
).order_by(HCRawRestingHeartRate.time.desc()).first()
|
|
|
|
return record.bpm if record else None
|
|
|
|
def get_daily_sleep_stats(self, user_id: int, target_date: date) -> Dict[str, int]:
|
|
"""
|
|
Get sleep statistics for a day.
|
|
|
|
Sleep is attributed to the day you wake up (e.g., sleep from 11pm Jan 7
|
|
to 7am Jan 8 is counted for Jan 8).
|
|
"""
|
|
start_of_day = datetime.combine(target_date, datetime.min.time())
|
|
end_of_day = datetime.combine(target_date + timedelta(days=1), datetime.min.time())
|
|
|
|
# Find sleep sessions that ended on this day
|
|
sessions = self.db.query(HCRawSleepSession).filter(
|
|
HCRawSleepSession.user_id == user_id,
|
|
HCRawSleepSession.end_time >= start_of_day,
|
|
HCRawSleepSession.end_time < end_of_day
|
|
).all()
|
|
|
|
total_minutes = 0
|
|
deep_minutes = 0
|
|
light_minutes = 0
|
|
rem_minutes = 0
|
|
awake_minutes = 0
|
|
|
|
for session in sessions:
|
|
# Total sleep duration
|
|
if session.start_time and session.end_time:
|
|
duration = (session.end_time - session.start_time).total_seconds() / 60
|
|
total_minutes += int(duration)
|
|
|
|
# Stage breakdown
|
|
for stage in session.stages:
|
|
if stage.start_time and stage.end_time:
|
|
stage_duration = (stage.end_time - stage.start_time).total_seconds() / 60
|
|
stage_minutes = int(stage_duration)
|
|
|
|
if stage.stage_name in ("deep", "STAGE_TYPE_DEEP"):
|
|
deep_minutes += stage_minutes
|
|
elif stage.stage_name in ("light", "STAGE_TYPE_LIGHT"):
|
|
light_minutes += stage_minutes
|
|
elif stage.stage_name in ("rem", "STAGE_TYPE_REM"):
|
|
rem_minutes += stage_minutes
|
|
elif stage.stage_name in ("awake", "awake_in_bed", "STAGE_TYPE_AWAKE"):
|
|
awake_minutes += stage_minutes
|
|
|
|
return {
|
|
"total": total_minutes,
|
|
"deep": deep_minutes,
|
|
"light": light_minutes,
|
|
"rem": rem_minutes,
|
|
"awake": awake_minutes
|
|
}
|
|
|
|
def get_daily_spo2_stats(self, user_id: int, target_date: date) -> Dict[str, Optional[float]]:
|
|
"""Get SpO2 statistics for a day."""
|
|
start_of_day = datetime.combine(target_date, datetime.min.time())
|
|
end_of_day = datetime.combine(target_date + timedelta(days=1), datetime.min.time())
|
|
|
|
records = self.db.query(HCRawOxygenSaturation).filter(
|
|
HCRawOxygenSaturation.user_id == user_id,
|
|
HCRawOxygenSaturation.time >= start_of_day,
|
|
HCRawOxygenSaturation.time < end_of_day
|
|
).all()
|
|
|
|
if not records:
|
|
return {"avg": None, "min": None, "max": None}
|
|
|
|
values = [r.percentage for r in records if r.percentage]
|
|
if not values:
|
|
return {"avg": None, "min": None, "max": None}
|
|
|
|
return {
|
|
"avg": sum(values) / len(values),
|
|
"min": min(values),
|
|
"max": max(values)
|
|
}
|
|
|
|
def get_daily_weight(self, user_id: int, target_date: date) -> Optional[float]:
|
|
"""Get the latest weight recorded on a day."""
|
|
start_of_day = datetime.combine(target_date, datetime.min.time())
|
|
end_of_day = datetime.combine(target_date + timedelta(days=1), datetime.min.time())
|
|
|
|
record = self.db.query(HCRawWeight).filter(
|
|
HCRawWeight.user_id == user_id,
|
|
HCRawWeight.time >= start_of_day,
|
|
HCRawWeight.time < end_of_day
|
|
).order_by(HCRawWeight.time.desc()).first()
|
|
|
|
return record.weight_kg if record else None
|
|
|
|
def get_daily_exercise_count(self, user_id: int, target_date: date) -> int:
|
|
"""Count exercise sessions for a day."""
|
|
start_of_day = datetime.combine(target_date, datetime.min.time())
|
|
end_of_day = datetime.combine(target_date + timedelta(days=1), datetime.min.time())
|
|
|
|
return self.db.query(func.count(HCRawExerciseSession.id)).filter(
|
|
HCRawExerciseSession.user_id == user_id,
|
|
HCRawExerciseSession.start_time >= start_of_day,
|
|
HCRawExerciseSession.start_time < end_of_day
|
|
).scalar() or 0
|
|
|
|
# ==================== COMBINED DAILY METRICS ====================
|
|
|
|
def get_daily_metrics(self, user_id: int, target_date: date) -> DailyMetrics:
|
|
"""Get all aggregated metrics for a single day."""
|
|
metrics = DailyMetrics(date=target_date)
|
|
|
|
# Steps and Activity
|
|
metrics.step_count = self.get_daily_steps(user_id, target_date)
|
|
metrics.distance_meters = self.get_daily_distance(user_id, target_date)
|
|
|
|
calories = self.get_daily_calories(user_id, target_date)
|
|
metrics.calories_burned = calories["total"]
|
|
metrics.active_calories = calories["active"]
|
|
|
|
# Heart Rate
|
|
hr_stats = self.get_daily_heart_rate_stats(user_id, target_date)
|
|
metrics.avg_heart_rate = hr_stats["avg"]
|
|
metrics.min_heart_rate = hr_stats["min"]
|
|
metrics.max_heart_rate = hr_stats["max"]
|
|
metrics.resting_heart_rate = self.get_daily_resting_heart_rate(user_id, target_date)
|
|
|
|
# Sleep
|
|
sleep_stats = self.get_daily_sleep_stats(user_id, target_date)
|
|
metrics.sleep_duration_minutes = sleep_stats["total"]
|
|
metrics.deep_sleep_minutes = sleep_stats["deep"]
|
|
metrics.light_sleep_minutes = sleep_stats["light"]
|
|
metrics.rem_sleep_minutes = sleep_stats["rem"]
|
|
metrics.awake_minutes = sleep_stats["awake"]
|
|
|
|
# SpO2
|
|
spo2_stats = self.get_daily_spo2_stats(user_id, target_date)
|
|
metrics.avg_spo2 = spo2_stats["avg"]
|
|
metrics.min_spo2 = spo2_stats["min"]
|
|
metrics.max_spo2 = spo2_stats["max"]
|
|
|
|
# Body
|
|
metrics.weight = self.get_daily_weight(user_id, target_date)
|
|
|
|
# Exercise
|
|
metrics.exercise_count = self.get_daily_exercise_count(user_id, target_date)
|
|
|
|
return metrics
|
|
|
|
# ==================== RANGE AGGREGATIONS ====================
|
|
|
|
def get_metrics_for_range(
|
|
self,
|
|
user_id: int,
|
|
start_date: date,
|
|
end_date: date
|
|
) -> List[DailyMetrics]:
|
|
"""Get daily metrics for a date range."""
|
|
metrics = []
|
|
current_date = start_date
|
|
|
|
while current_date <= end_date:
|
|
daily = self.get_daily_metrics(user_id, current_date)
|
|
metrics.append(daily)
|
|
current_date += timedelta(days=1)
|
|
|
|
return metrics
|
|
|
|
def get_weekly_summary(self, user_id: int, ref_date: date) -> Dict[str, Any]:
|
|
"""Get aggregated summary for a week (7 days ending on ref_date)."""
|
|
start_date = ref_date - timedelta(days=6)
|
|
daily_metrics = self.get_metrics_for_range(user_id, start_date, ref_date)
|
|
|
|
# Aggregate
|
|
total_steps = sum(m.step_count for m in daily_metrics)
|
|
total_distance = sum(m.distance_meters for m in daily_metrics)
|
|
total_calories = sum(m.calories_burned for m in daily_metrics)
|
|
|
|
# Average heart rate (excluding days with no data)
|
|
hr_values = [m.avg_heart_rate for m in daily_metrics if m.avg_heart_rate]
|
|
avg_hr = int(sum(hr_values) / len(hr_values)) if hr_values else None
|
|
|
|
# Average sleep
|
|
sleep_values = [m.sleep_duration_minutes for m in daily_metrics if m.sleep_duration_minutes > 0]
|
|
avg_sleep = int(sum(sleep_values) / len(sleep_values)) if sleep_values else 0
|
|
|
|
# Latest weight
|
|
weights = [m.weight for m in daily_metrics if m.weight]
|
|
latest_weight = weights[-1] if weights else None
|
|
|
|
return {
|
|
"period": "week",
|
|
"start_date": start_date.isoformat(),
|
|
"end_date": ref_date.isoformat(),
|
|
"days_with_data": len([m for m in daily_metrics if m.step_count > 0]),
|
|
"totals": {
|
|
"steps": total_steps,
|
|
"distance_meters": total_distance,
|
|
"calories_burned": total_calories,
|
|
"exercise_count": sum(m.exercise_count for m in daily_metrics),
|
|
},
|
|
"averages": {
|
|
"steps_per_day": total_steps // 7,
|
|
"distance_per_day": total_distance / 7,
|
|
"calories_per_day": total_calories / 7,
|
|
"sleep_minutes": avg_sleep,
|
|
"heart_rate": avg_hr,
|
|
},
|
|
"latest": {
|
|
"weight": latest_weight,
|
|
},
|
|
"daily": [m.to_dict() for m in daily_metrics]
|
|
}
|
|
|
|
def get_monthly_summary(self, user_id: int, ref_date: date) -> Dict[str, Any]:
|
|
"""Get aggregated summary for a month (30 days ending on ref_date)."""
|
|
start_date = ref_date - timedelta(days=29)
|
|
daily_metrics = self.get_metrics_for_range(user_id, start_date, ref_date)
|
|
|
|
# Aggregate
|
|
total_steps = sum(m.step_count for m in daily_metrics)
|
|
total_distance = sum(m.distance_meters for m in daily_metrics)
|
|
total_calories = sum(m.calories_burned for m in daily_metrics)
|
|
|
|
days_count = len(daily_metrics)
|
|
|
|
return {
|
|
"period": "month",
|
|
"start_date": start_date.isoformat(),
|
|
"end_date": ref_date.isoformat(),
|
|
"days_with_data": len([m for m in daily_metrics if m.step_count > 0]),
|
|
"totals": {
|
|
"steps": total_steps,
|
|
"distance_meters": total_distance,
|
|
"calories_burned": total_calories,
|
|
"exercise_count": sum(m.exercise_count for m in daily_metrics),
|
|
},
|
|
"averages": {
|
|
"steps_per_day": total_steps // days_count if days_count > 0 else 0,
|
|
"distance_per_day": total_distance / days_count if days_count > 0 else 0,
|
|
"calories_per_day": total_calories / days_count if days_count > 0 else 0,
|
|
},
|
|
}
|
|
|
|
# ==================== TODAY'S METRICS (CURRENT/LIVE) ====================
|
|
|
|
def get_today_metrics(self, user_id: int) -> DailyMetrics:
|
|
"""
|
|
Get today's metrics - useful for dashboard display.
|
|
This returns the most up-to-date aggregation for the current day.
|
|
"""
|
|
return self.get_daily_metrics(user_id, date.today())
|
|
|
|
# ==================== METRIC HISTORY ====================
|
|
|
|
def get_metric_history(
|
|
self,
|
|
user_id: int,
|
|
metric_type: str,
|
|
start_date: date,
|
|
end_date: date
|
|
) -> List[Dict[str, Any]]:
|
|
"""
|
|
Get historical data for a specific metric type.
|
|
|
|
Args:
|
|
metric_type: One of "steps", "heart_rate", "sleep", "weight", "spo2"
|
|
"""
|
|
history = []
|
|
current_date = start_date
|
|
|
|
while current_date <= end_date:
|
|
data_point = {"date": current_date.isoformat()}
|
|
|
|
if metric_type == "steps":
|
|
data_point["value"] = self.get_daily_steps(user_id, current_date)
|
|
elif metric_type == "distance":
|
|
data_point["value"] = self.get_daily_distance(user_id, current_date)
|
|
elif metric_type == "calories":
|
|
cals = self.get_daily_calories(user_id, current_date)
|
|
data_point["value"] = cals["total"]
|
|
data_point["active"] = cals["active"]
|
|
elif metric_type == "heart_rate":
|
|
hr = self.get_daily_heart_rate_stats(user_id, current_date)
|
|
data_point["avg"] = hr["avg"]
|
|
data_point["min"] = hr["min"]
|
|
data_point["max"] = hr["max"]
|
|
data_point["resting"] = self.get_daily_resting_heart_rate(user_id, current_date)
|
|
elif metric_type == "sleep":
|
|
sleep = self.get_daily_sleep_stats(user_id, current_date)
|
|
data_point["total"] = sleep["total"]
|
|
data_point["deep"] = sleep["deep"]
|
|
data_point["light"] = sleep["light"]
|
|
data_point["rem"] = sleep["rem"]
|
|
elif metric_type == "weight":
|
|
data_point["value"] = self.get_daily_weight(user_id, current_date)
|
|
elif metric_type == "spo2":
|
|
spo2 = self.get_daily_spo2_stats(user_id, current_date)
|
|
data_point["avg"] = spo2["avg"]
|
|
data_point["min"] = spo2["min"]
|
|
data_point["max"] = spo2["max"]
|
|
|
|
history.append(data_point)
|
|
current_date += timedelta(days=1)
|
|
|
|
return history
|