Files
HabitForge/backend/api/v1/health_connect_raw.py
T
2026-07-30 23:25:20 -04:00

735 lines
26 KiB
Python

"""
Health Connect Raw Data API - Endpoints for receiving raw health records from Android.
This endpoint receives individual health records (not aggregated) from the Android
companion app and stores them for accurate on-demand aggregation.
OPTIMIZED VERSION: Uses batch inserts and efficient duplicate detection.
"""
from fastapi import APIRouter, Depends, HTTPException, status, BackgroundTasks
from sqlalchemy.orm import Session
from sqlalchemy import and_, tuple_
from typing import List, Optional, Set, Tuple
from pydantic import BaseModel
from datetime import datetime
import logging
import time
logger = logging.getLogger(__name__)
from backend.database import get_db
from backend.models.user import User
from backend.models.health_connect import HealthConnectSyncLog
from backend.models.health_connect_raw import (
HCRawSteps,
HCRawHeartRate,
HCRawHeartRateSample,
HCRawRestingHeartRate,
HCRawSleepSession,
HCRawSleepStage,
HCRawDistance,
HCRawCalories,
HCRawOxygenSaturation,
HCRawWeight,
HCRawHeight,
HCRawBodyFat,
HCRawExerciseSession,
HCRawNutrition,
HCRawHydration,
SLEEP_STAGE_TYPES,
)
from backend import auth
router = APIRouter(prefix="/api/health-connect", tags=["health-connect-raw"])
# ==================== REQUEST MODELS ====================
class RawStepsPayload(BaseModel):
"""Individual step record from Health Connect"""
record_id: str
data_source: str
data_source_name: Optional[str] = None
start_time: datetime
end_time: datetime
count: int
class HeartRateSamplePayload(BaseModel):
"""Single heart rate sample"""
time: datetime
bpm: int
class RawHeartRatePayload(BaseModel):
"""Heart rate record with samples"""
record_id: str
data_source: str
data_source_name: Optional[str] = None
start_time: datetime
end_time: datetime
samples: List[HeartRateSamplePayload] = []
class RawRestingHeartRatePayload(BaseModel):
"""Resting heart rate record"""
record_id: str
data_source: str
data_source_name: Optional[str] = None
time: datetime
bpm: int
class SleepStagePayload(BaseModel):
"""Sleep stage within a session"""
start_time: datetime
end_time: datetime
stage_type: int
stage_name: Optional[str] = None
class RawSleepSessionPayload(BaseModel):
"""Sleep session with stages"""
record_id: str
data_source: str
data_source_name: Optional[str] = None
start_time: datetime
end_time: datetime
title: Optional[str] = None
notes: Optional[str] = None
stages: List[SleepStagePayload] = []
class RawDistancePayload(BaseModel):
"""Distance record"""
record_id: str
data_source: str
data_source_name: Optional[str] = None
start_time: datetime
end_time: datetime
distance_meters: float
class RawCaloriesPayload(BaseModel):
"""Calories record (active or total)"""
record_id: str
data_source: str
data_source_name: Optional[str] = None
start_time: datetime
end_time: datetime
calories: float
record_type: str # "active" or "total"
class RawOxygenSaturationPayload(BaseModel):
"""SpO2 record"""
record_id: str
data_source: str
data_source_name: Optional[str] = None
time: datetime
percentage: float
class RawWeightPayload(BaseModel):
"""Weight record"""
record_id: str
data_source: str
data_source_name: Optional[str] = None
time: datetime
weight_kg: float
class RawHeightPayload(BaseModel):
"""Height record"""
record_id: str
data_source: str
data_source_name: Optional[str] = None
time: datetime
height_meters: float
class RawBodyFatPayload(BaseModel):
"""Body fat record"""
record_id: str
data_source: str
data_source_name: Optional[str] = None
time: datetime
percentage: float
class RawExerciseSessionPayload(BaseModel):
"""Exercise session"""
record_id: str
data_source: str
data_source_name: Optional[str] = None
start_time: datetime
end_time: datetime
exercise_type: int
exercise_type_name: Optional[str] = None
title: Optional[str] = None
notes: Optional[str] = None
distance_meters: Optional[float] = None
calories: Optional[float] = None
avg_heart_rate: Optional[int] = None
max_heart_rate: Optional[int] = None
min_heart_rate: Optional[int] = None
steps: Optional[int] = None
elevation_gained_meters: Optional[float] = None
class RawNutritionPayload(BaseModel):
"""Nutrition record"""
record_id: str
data_source: str
data_source_name: Optional[str] = None
start_time: datetime
end_time: datetime
name: Optional[str] = None
meal_type: Optional[int] = None
calories: Optional[float] = None
protein_grams: Optional[float] = None
carbohydrates_grams: Optional[float] = None
fat_grams: Optional[float] = None
fiber_grams: Optional[float] = None
sugar_grams: Optional[float] = None
class RawHydrationPayload(BaseModel):
"""Hydration record"""
record_id: str
data_source: str
data_source_name: Optional[str] = None
start_time: datetime
end_time: datetime
volume_liters: float
class RawDataSyncRequest(BaseModel):
"""Complete raw data sync request from Android app"""
device_id: str
sync_timestamp: str = ""
# All record types
steps_records: List[RawStepsPayload] = []
heart_rate_records: List[RawHeartRatePayload] = []
resting_hr_records: List[RawRestingHeartRatePayload] = []
sleep_sessions: List[RawSleepSessionPayload] = []
distance_records: List[RawDistancePayload] = []
calories_records: List[RawCaloriesPayload] = []
spo2_records: List[RawOxygenSaturationPayload] = []
weight_records: List[RawWeightPayload] = []
height_records: List[RawHeightPayload] = []
body_fat_records: List[RawBodyFatPayload] = []
exercise_sessions: List[RawExerciseSessionPayload] = []
nutrition_records: List[RawNutritionPayload] = []
hydration_records: List[RawHydrationPayload] = []
class RawSyncResponse(BaseModel):
"""Response for raw data sync"""
status: str
sync_id: int
records_synced: dict # Count per record type
message: str
# ==================== OPTIMIZED HELPER FUNCTIONS ====================
def get_existing_record_ids(db: Session, model, user_id: int, record_ids: List[str]) -> Set[str]:
"""
Batch query to find which record_ids already exist for the user.
Returns a set of existing record_ids.
"""
if not record_ids:
return set()
# Process in batches to avoid query limits
batch_size = 500
existing_ids = set()
for i in range(0, len(record_ids), batch_size):
batch = record_ids[i:i + batch_size]
results = db.query(model.record_id).filter(
model.user_id == user_id,
model.record_id.in_(batch)
).all()
existing_ids.update(r[0] for r in results)
return existing_ids
def batch_insert_simple_records(db: Session, model, user_id: int, records: list,
sync_id: int, existing_ids: Set[str]) -> int:
"""
Batch insert records, skipping those that already exist.
Also deduplicates records within the payload itself.
Returns count of new records inserted.
"""
new_records = []
now = datetime.utcnow()
seen_ids = set() # Track record_ids we've already processed in this batch
for record in records:
# Skip if already exists in database
if record.record_id in existing_ids:
continue
# Skip if we've already seen this record_id in this payload (dedupe)
if record.record_id in seen_ids:
continue
seen_ids.add(record.record_id)
data = record.model_dump()
data['user_id'] = user_id
data['sync_id'] = sync_id
data['synced_at'] = now
new_records.append(model(**data))
if new_records:
db.bulk_save_objects(new_records)
return len(new_records)
# ==================== ENDPOINTS ====================
@router.post("/sync-raw", response_model=RawSyncResponse, status_code=status.HTTP_200_OK)
async def receive_raw_health_data(
request: RawDataSyncRequest,
db: Session = Depends(get_db),
current_user: User = Depends(auth.get_current_user)
):
"""
Receive raw health records from Android companion app.
OPTIMIZED: Uses batch queries and bulk inserts for much better performance.
"""
start_time = time.time()
# Count total records (excluding HR samples for now)
total_records = (
len(request.steps_records) +
len(request.heart_rate_records) +
len(request.resting_hr_records) +
len(request.sleep_sessions) +
len(request.distance_records) +
len(request.calories_records) +
len(request.spo2_records) +
len(request.weight_records) +
len(request.height_records) +
len(request.body_fat_records) +
len(request.exercise_sessions) +
len(request.nutrition_records) +
len(request.hydration_records)
)
# Count HR samples
total_hr_samples = sum(len(hr.samples) for hr in request.heart_rate_records)
logger.info(f"=== Raw Data Sync from user {current_user.username} ===")
logger.info(f"Device: {request.device_id}, Records: {total_records}, HR Samples: {total_hr_samples}")
# Create sync log
sync_log = HealthConnectSyncLog(
user_id=current_user.id,
device_id=request.device_id,
records_pushed=total_records,
status="processing"
)
db.add(sync_log)
db.flush()
# Track counts
counts = {
"steps": 0,
"heart_rate": 0,
"heart_rate_samples": 0,
"resting_hr": 0,
"sleep_sessions": 0,
"sleep_stages": 0,
"distance": 0,
"calories": 0,
"spo2": 0,
"weight": 0,
"height": 0,
"body_fat": 0,
"exercise": 0,
"nutrition": 0,
"hydration": 0,
}
try:
now = datetime.utcnow()
# === STEPS ===
if request.steps_records:
step_ids = [r.record_id for r in request.steps_records]
existing_step_ids = get_existing_record_ids(db, HCRawSteps, current_user.id, step_ids)
counts["steps"] = batch_insert_simple_records(
db, HCRawSteps, current_user.id, request.steps_records,
sync_log.id, existing_step_ids
)
logger.info(f"Steps: {counts['steps']} new of {len(request.steps_records)} total")
# === DISTANCE ===
if request.distance_records:
dist_ids = [r.record_id for r in request.distance_records]
existing_dist_ids = get_existing_record_ids(db, HCRawDistance, current_user.id, dist_ids)
counts["distance"] = batch_insert_simple_records(
db, HCRawDistance, current_user.id, request.distance_records,
sync_log.id, existing_dist_ids
)
logger.info(f"Distance: {counts['distance']} new of {len(request.distance_records)} total")
# === CALORIES ===
if request.calories_records:
cal_ids = [r.record_id for r in request.calories_records]
existing_cal_ids = get_existing_record_ids(db, HCRawCalories, current_user.id, cal_ids)
counts["calories"] = batch_insert_simple_records(
db, HCRawCalories, current_user.id, request.calories_records,
sync_log.id, existing_cal_ids
)
logger.info(f"Calories: {counts['calories']} new of {len(request.calories_records)} total")
# === RESTING HR ===
if request.resting_hr_records:
rhr_ids = [r.record_id for r in request.resting_hr_records]
existing_rhr_ids = get_existing_record_ids(db, HCRawRestingHeartRate, current_user.id, rhr_ids)
counts["resting_hr"] = batch_insert_simple_records(
db, HCRawRestingHeartRate, current_user.id, request.resting_hr_records,
sync_log.id, existing_rhr_ids
)
logger.info(f"Resting HR: {counts['resting_hr']} new of {len(request.resting_hr_records)} total")
# === SpO2 ===
if request.spo2_records:
spo2_ids = [r.record_id for r in request.spo2_records]
existing_spo2_ids = get_existing_record_ids(db, HCRawOxygenSaturation, current_user.id, spo2_ids)
counts["spo2"] = batch_insert_simple_records(
db, HCRawOxygenSaturation, current_user.id, request.spo2_records,
sync_log.id, existing_spo2_ids
)
logger.info(f"SpO2: {counts['spo2']} new of {len(request.spo2_records)} total")
# === WEIGHT ===
if request.weight_records:
weight_ids = [r.record_id for r in request.weight_records]
existing_weight_ids = get_existing_record_ids(db, HCRawWeight, current_user.id, weight_ids)
counts["weight"] = batch_insert_simple_records(
db, HCRawWeight, current_user.id, request.weight_records,
sync_log.id, existing_weight_ids
)
logger.info(f"Weight: {counts['weight']} new of {len(request.weight_records)} total")
# === HEIGHT ===
if request.height_records:
height_ids = [r.record_id for r in request.height_records]
existing_height_ids = get_existing_record_ids(db, HCRawHeight, current_user.id, height_ids)
counts["height"] = batch_insert_simple_records(
db, HCRawHeight, current_user.id, request.height_records,
sync_log.id, existing_height_ids
)
logger.info(f"Height: {counts['height']} new of {len(request.height_records)} total")
# === BODY FAT ===
if request.body_fat_records:
bf_ids = [r.record_id for r in request.body_fat_records]
existing_bf_ids = get_existing_record_ids(db, HCRawBodyFat, current_user.id, bf_ids)
counts["body_fat"] = batch_insert_simple_records(
db, HCRawBodyFat, current_user.id, request.body_fat_records,
sync_log.id, existing_bf_ids
)
logger.info(f"Body Fat: {counts['body_fat']} new of {len(request.body_fat_records)} total")
# === NUTRITION ===
if request.nutrition_records:
nut_ids = [r.record_id for r in request.nutrition_records]
existing_nut_ids = get_existing_record_ids(db, HCRawNutrition, current_user.id, nut_ids)
counts["nutrition"] = batch_insert_simple_records(
db, HCRawNutrition, current_user.id, request.nutrition_records,
sync_log.id, existing_nut_ids
)
logger.info(f"Nutrition: {counts['nutrition']} new of {len(request.nutrition_records)} total")
# === HYDRATION ===
if request.hydration_records:
hyd_ids = [r.record_id for r in request.hydration_records]
existing_hyd_ids = get_existing_record_ids(db, HCRawHydration, current_user.id, hyd_ids)
counts["hydration"] = batch_insert_simple_records(
db, HCRawHydration, current_user.id, request.hydration_records,
sync_log.id, existing_hyd_ids
)
logger.info(f"Hydration: {counts['hydration']} new of {len(request.hydration_records)} total")
# === EXERCISE SESSIONS ===
if request.exercise_sessions:
ex_ids = [r.record_id for r in request.exercise_sessions]
existing_ex_ids = get_existing_record_ids(db, HCRawExerciseSession, current_user.id, ex_ids)
counts["exercise"] = batch_insert_simple_records(
db, HCRawExerciseSession, current_user.id, request.exercise_sessions,
sync_log.id, existing_ex_ids
)
logger.info(f"Exercise: {counts['exercise']} new of {len(request.exercise_sessions)} total")
# === HEART RATE (with samples - needs special handling) ===
if request.heart_rate_records:
hr_ids = [r.record_id for r in request.heart_rate_records]
existing_hr_ids = get_existing_record_ids(db, HCRawHeartRate, current_user.id, hr_ids)
# Insert new HR records
new_hr_records = []
hr_records_for_samples = [] # Track which records need samples
for record in request.heart_rate_records:
if record.record_id in existing_hr_ids:
continue
data = record.model_dump(exclude={"samples"})
data['user_id'] = current_user.id
data['sync_id'] = sync_log.id
data['synced_at'] = now
hr_obj = HCRawHeartRate(**data)
new_hr_records.append(hr_obj)
hr_records_for_samples.append((hr_obj, record.samples))
if new_hr_records:
db.bulk_save_objects(new_hr_records, return_defaults=True)
db.flush() # Get IDs
counts["heart_rate"] = len(new_hr_records)
# Now insert samples in batches
all_samples = []
for hr_obj, samples in hr_records_for_samples:
for sample in samples:
all_samples.append(HCRawHeartRateSample(
record_id=hr_obj.id,
time=sample.time,
bpm=sample.bpm
))
if all_samples:
# Insert in batches of 1000
batch_size = 1000
for i in range(0, len(all_samples), batch_size):
batch = all_samples[i:i + batch_size]
db.bulk_save_objects(batch)
counts["heart_rate_samples"] = len(all_samples)
logger.info(f"Heart Rate: {counts['heart_rate']} records, {counts['heart_rate_samples']} samples")
# === SLEEP SESSIONS (with stages) ===
if request.sleep_sessions:
sleep_ids = [r.record_id for r in request.sleep_sessions]
existing_sleep_ids = get_existing_record_ids(db, HCRawSleepSession, current_user.id, sleep_ids)
new_sleep_records = []
sleep_records_for_stages = []
for record in request.sleep_sessions:
if record.record_id in existing_sleep_ids:
continue
data = record.model_dump(exclude={"stages"})
data['user_id'] = current_user.id
data['sync_id'] = sync_log.id
data['synced_at'] = now
sleep_obj = HCRawSleepSession(**data)
new_sleep_records.append(sleep_obj)
sleep_records_for_stages.append((sleep_obj, record.stages))
if new_sleep_records:
db.bulk_save_objects(new_sleep_records, return_defaults=True)
db.flush()
counts["sleep_sessions"] = len(new_sleep_records)
# Insert stages
all_stages = []
for sleep_obj, stages in sleep_records_for_stages:
for stage in stages:
stage_name = stage.stage_name or SLEEP_STAGE_TYPES.get(stage.stage_type, "unknown")
all_stages.append(HCRawSleepStage(
session_id=sleep_obj.id,
start_time=stage.start_time,
end_time=stage.end_time,
stage_type=stage.stage_type,
stage_name=stage_name
))
if all_stages:
db.bulk_save_objects(all_stages)
counts["sleep_stages"] = len(all_stages)
logger.info(f"Sleep: {counts['sleep_sessions']} sessions, {counts['sleep_stages']} stages")
# Update sync log
sync_log.status = "success"
sync_log.records_pushed = sum(counts.values())
db.commit()
elapsed = time.time() - start_time
logger.info(f"=== Raw Sync Complete in {elapsed:.2f}s ===")
logger.info(f"Records synced: {counts}")
return RawSyncResponse(
status="success",
sync_id=sync_log.id,
records_synced=counts,
message=f"Successfully synced {sum(counts.values())} raw records in {elapsed:.2f}s"
)
except Exception as e:
logger.error(f"Raw sync failed: {str(e)}", exc_info=True)
db.rollback()
sync_log.status = "error"
sync_log.error_message = str(e)[:500] # Limit error message length
db.commit()
raise HTTPException(
status_code=status.HTTP_500_INTERNAL_SERVER_ERROR,
detail=f"Raw sync failed: {str(e)}"
)
@router.get("/raw/steps", status_code=status.HTTP_200_OK)
async def get_raw_steps(
start_date: Optional[datetime] = None,
end_date: Optional[datetime] = None,
limit: int = 1000,
db: Session = Depends(get_db),
current_user: User = Depends(auth.get_current_user)
):
"""Get raw step records for a date range."""
from datetime import timedelta
if not end_date:
end_date = datetime.utcnow()
if not start_date:
start_date = end_date - timedelta(days=7)
records = db.query(HCRawSteps).filter(
HCRawSteps.user_id == current_user.id,
HCRawSteps.start_time >= start_date,
HCRawSteps.start_time <= end_date
).order_by(HCRawSteps.start_time.desc()).limit(limit).all()
return [
{
"id": r.id,
"record_id": r.record_id,
"data_source": r.data_source,
"data_source_name": r.data_source_name,
"start_time": r.start_time.isoformat() if r.start_time else None,
"end_time": r.end_time.isoformat() if r.end_time else None,
"count": r.count
}
for r in records
]
@router.get("/raw/heart-rate", status_code=status.HTTP_200_OK)
async def get_raw_heart_rate(
start_date: Optional[datetime] = None,
end_date: Optional[datetime] = None,
include_samples: bool = False,
limit: int = 100,
db: Session = Depends(get_db),
current_user: User = Depends(auth.get_current_user)
):
"""Get raw heart rate records for a date range."""
from datetime import timedelta
if not end_date:
end_date = datetime.utcnow()
if not start_date:
start_date = end_date - timedelta(days=1)
records = db.query(HCRawHeartRate).filter(
HCRawHeartRate.user_id == current_user.id,
HCRawHeartRate.start_time >= start_date,
HCRawHeartRate.start_time <= end_date
).order_by(HCRawHeartRate.start_time.desc()).limit(limit).all()
result = []
for r in records:
record_dict = {
"id": r.id,
"record_id": r.record_id,
"data_source": r.data_source,
"start_time": r.start_time.isoformat() if r.start_time else None,
"end_time": r.end_time.isoformat() if r.end_time else None,
}
if include_samples:
record_dict["samples"] = [
{"time": s.time.isoformat(), "bpm": s.bpm}
for s in r.samples
]
else:
record_dict["sample_count"] = len(r.samples)
result.append(record_dict)
return result
@router.get("/raw/stats", status_code=status.HTTP_200_OK)
async def get_raw_data_stats(
db: Session = Depends(get_db),
current_user: User = Depends(auth.get_current_user)
):
"""Get statistics about raw data stored for current user."""
from sqlalchemy import func
stats = {
"steps": db.query(func.count(HCRawSteps.id)).filter(
HCRawSteps.user_id == current_user.id
).scalar(),
"heart_rate": db.query(func.count(HCRawHeartRate.id)).filter(
HCRawHeartRate.user_id == current_user.id
).scalar(),
"heart_rate_samples": db.query(func.count(HCRawHeartRateSample.id)).join(
HCRawHeartRate
).filter(
HCRawHeartRate.user_id == current_user.id
).scalar(),
"resting_hr": db.query(func.count(HCRawRestingHeartRate.id)).filter(
HCRawRestingHeartRate.user_id == current_user.id
).scalar(),
"sleep_sessions": db.query(func.count(HCRawSleepSession.id)).filter(
HCRawSleepSession.user_id == current_user.id
).scalar(),
"distance": db.query(func.count(HCRawDistance.id)).filter(
HCRawDistance.user_id == current_user.id
).scalar(),
"calories": db.query(func.count(HCRawCalories.id)).filter(
HCRawCalories.user_id == current_user.id
).scalar(),
"spo2": db.query(func.count(HCRawOxygenSaturation.id)).filter(
HCRawOxygenSaturation.user_id == current_user.id
).scalar(),
"weight": db.query(func.count(HCRawWeight.id)).filter(
HCRawWeight.user_id == current_user.id
).scalar(),
"exercise": db.query(func.count(HCRawExerciseSession.id)).filter(
HCRawExerciseSession.user_id == current_user.id
).scalar(),
}
# Get date range
oldest_step = db.query(func.min(HCRawSteps.start_time)).filter(
HCRawSteps.user_id == current_user.id
).scalar()
newest_step = db.query(func.max(HCRawSteps.start_time)).filter(
HCRawSteps.user_id == current_user.id
).scalar()
return {
"record_counts": stats,
"total_records": sum(stats.values()),
"date_range": {
"oldest": oldest_step.isoformat() if oldest_step else None,
"newest": newest_step.isoformat() if newest_step else None
}
}