Skip to main content
Synced from the repo — do not edit here

Canonical source: UNIFIED_DATA_ARCHITECTURE_GUIDE.md. This page is generated by docs/scripts/sync-handbook.mjs. Edit the source file in the repo; changes appear here on the next build.

Unified Data Architecture Implementation Guide

Overview

This guide documents the implementation of the unified data architecture that enables seamless querying across Apple Health, Terra API, and Basis Events data sources. All data is normalized to the BasisSummaryV1 format for unified storage and cross-source correlation.

Architecture Summary

graph TD
A[Apple Health] --> B[convertHealthPointsToSummaries]
B --> C[BasisSummaryV1]

D[Terra API] --> E[TerraToBasisisConverter]
E --> C

F[Basis Events] --> G[BasisEventToSummaryConverter]
G --> C

C --> H[UnifiedDataManager]
H --> I[DuckDB Storage]

J[AI Copilot] --> K[functions_ai_unified]
K --> H
K --> I

L[Clinician Query: "George's max HR May 2025"] --> K
M[Fasting Analysis: "HR + Glucose during fast"] --> K

Implementation Files

1. Data Converters

terra_to_basis_converter.py - Converts Terra webhook data to BasisSummaryV1

  • Heart rate, glucose, steps, HRV, body measurements
  • Device information mapping
  • Timestamp and timezone handling

basis_event_to_summary_converter.py - Converts Basis Events to BasisSummaryV1

  • Fasting timers, workouts, medication events
  • Duration calculation and event context
  • Timer and user event support

2. Unified Storage

unified_data_manager.py - Central manager for all data sources

  • Cross-source storage and querying
  • Device deduplication and registry
  • Performance optimized DuckDB operations
  • Specialized queries (fasting analysis, max HR, etc.)

3. Enhanced AI Functions

functions_ai_unified.py - AI tools for cross-source querying

  • query_unified_health_data: Query across all sources
  • analyze_fasting_periods: Correlate events with health data
  • get_max_heart_rate_cross_device: Multi-device aggregation
  • Enhanced LangChain agent with unified tools

Integration Steps

Step 1: Update Terra Webhook Processing

Modify your existing Terra webhook handler to use the unified converter:

# In your existing terra_adapter.py or webhook handler
from src.terra_to_basis_converter import convert_terra_to_basis_summaries
from src.unified_data_manager import get_unified_data_manager

def process_terra_webhook_unified(request_type: str, json_data: dict, user_id: str) -> bool:
try:
# Get unified data manager
data_manager = get_unified_data_manager(user_id)

# Store Terra data in unified format
success = data_manager.store_terra_data(json_data, request_type)

if success:
# Sync to GCS as before
db_manager = get_database_manager()
db_manager.sync_to_gcs(user_id)

return success
except Exception as e:
logger.exception(f"Error processing Terra webhook for user {user_id}: {e}")
return False

Step 2: Update Apple Health Processing

Modify your existing Apple Health sync to use unified storage:

# In your existing health sync service
from src.unified_data_manager import get_unified_data_manager

def sync_apple_health_unified(user_id: str, summaries: List[BasisSummaryV1]) -> bool:
try:
# Get unified data manager
data_manager = get_unified_data_manager(user_id)

# Store Apple Health data (already in BasisSummaryV1 format)
success = data_manager.store_apple_health_summaries(summaries)

return success
except Exception as e:
logger.exception(f"Error syncing Apple Health for user {user_id}: {e}")
return False

Step 3: Update Basis Events Processing

Add unified storage for Basis Events:

# In your existing event processing service
from src.unified_data_manager import get_unified_data_manager

def store_basis_events_unified(user_id: str, events: List[Dict[str, Any]]) -> bool:
try:
# Get unified data manager
data_manager = get_unified_data_manager(user_id)

# Store Basis Events in unified format
success = data_manager.store_basis_events(events)

return success
except Exception as e:
logger.exception(f"Error storing Basis events for user {user_id}: {e}")
return False

Step 4: Update AI Copilot

Replace or enhance your existing AI functions:

# In your AI copilot service
from src.functions_ai_unified import (
create_unified_clinical_agent,
query_unified_health_data,
analyze_fasting_periods,
get_max_heart_rate_cross_device
)

def create_enhanced_copilot():
"""Create AI copilot with unified data access."""
return create_unified_clinical_agent()

# Example usage in your copilot endpoint
async def handle_copilot_query(user_id: str, query: str):
agent = create_enhanced_copilot()

# The agent can now handle cross-source queries like:
# - "What was my max heart rate in May 2025?"
# - "Show me glucose and heart rate during my last fasting period"
# - "Compare my Oura Ring and Apple Watch step counts"

response = await agent.arun(f"User {user_id}: {query}")
return response

Example Use Cases

1. Cross-Device Heart Rate Query

Query: "What was George's max heart rate in May 2025?"

Data Sources: Apple Watch (Apple Health) + Oura Ring (Terra)

Implementation:

query = UnifiedHealthQuery(
user_id="george_123",
health_types=["heart_rate"],
start_date="2025-05-01",
end_date="2025-05-31",
aggregation="max"
)

result = query_unified_health_data(query)
# Returns max HR across all devices with device attribution

2. Fasting Event Analysis

Query: "Show me heart rate and glucose during my last fasting period"

Data Sources: Fasting Timer (Basis Events) + Heart Rate (Apple/Terra) + Glucose (Terra CGM)

Implementation:

query = FastingAnalysisQuery(
user_id="user_123",
start_date="2025-01-01",
end_date="2025-01-31",
include_metrics=["heart_rate", "glucose", "hrv"]
)

result = analyze_fasting_periods(query)
# Returns fasting periods with correlated physiological data

3. Multi-Source Event Summary

Query: "Create a workout summary with heart rate data"

Data Sources: Workout Event (Basis) + Heart Rate (Apple Health/Terra)

Process:

  1. Workout event converted to BasisSummaryV1 with duration
  2. Heart rate data queried for workout time period
  3. Cross-correlation provides workout with HR zones

Database Schema Changes

The unified architecture uses your existing DuckDB schema but adds:

  1. Source Attribution: summary_id field with prefixes:

    • u_xxx: Unified summaries from converters
    • a_xxx: Apple Health summaries
    • t_xxx: Terra summaries
    • b_xxx: Basis event summaries
  2. Device Registry: Enhanced device management across sources

  3. Cross-Reference Tables: Link events to physiological data

Performance Considerations

1. Data Volume

  • Each data source generates different volumes
  • Terra: High-frequency (1-minute intervals)
  • Apple Health: Variable frequency
  • Basis Events: Low frequency, high importance

2. Query Optimization

  • Time-based indexing for correlation queries
  • Device-based partitioning for multi-device scenarios
  • Materialized views for common aggregations

3. Storage Efficiency

  • Unified format reduces storage overhead
  • Device deduplication prevents redundancy
  • Compression for historical data

Testing Strategy

1. Unit Tests

# Test converters
def test_terra_converter():
terra_data = load_sample_terra_webhook()
summaries = convert_terra_to_basis_summaries("user123", terra_data, "body")
assert len(summaries) > 0
assert all(s.type in [BasisHealthType.HEART_RATE, BasisHealthType.GLUCOSE] for s in summaries)

def test_basis_event_converter():
fasting_event = create_sample_fasting_event()
summary = convert_fasting_event_to_summary(fasting_event)
assert summary.type == BasisHealthType.FASTING
assert summary.vt == BasisValueType.FAST

2. Integration Tests

def test_cross_source_query():
# Setup test data from all sources
setup_test_apple_health_data("user123")
setup_test_terra_data("user123")
setup_test_basis_events("user123")

# Query across all sources
data_manager = get_unified_data_manager("user123")
results = data_manager.query_cross_source_data(
health_types=[BasisHealthType.HEART_RATE],
start_time=datetime(2025, 1, 1),
end_time=datetime(2025, 1, 31)
)

# Verify cross-source results
apple_results = [r for r in results if r['summary_id'].startswith('a_')]
terra_results = [r for r in results if r['summary_id'].startswith('t_')]

assert len(apple_results) > 0
assert len(terra_results) > 0

3. AI Function Tests

def test_unified_ai_functions():
# Test max heart rate query
result = get_max_heart_rate_cross_device("user123", "2025-05-01", "2025-05-31")
data = json.loads(result)
assert data['success'] == True
assert 'max_heart_rate' in data

# Test fasting analysis
query = FastingAnalysisQuery(
user_id="user123",
start_date="2025-01-01",
end_date="2025-01-31"
)
result = analyze_fasting_periods(query)
data = json.loads(result)
assert data['success'] == True
assert 'fasting_periods_analyzed' in data

Migration Strategy

Phase 1: Parallel Implementation

  • Deploy unified converters alongside existing system
  • Process new data through both old and new pipelines
  • Verify data consistency

Phase 2: Query Migration

  • Update AI functions to use unified queries
  • Maintain backward compatibility
  • Monitor performance impact

Phase 3: Complete Migration

  • Switch all data processing to unified pipeline
  • Remove legacy data processing code
  • Optimize performance based on usage patterns

Monitoring and Alerting

Data Quality Metrics

  • Conversion success rates by source
  • Data volume consistency checks
  • Cross-source correlation validation

Performance Metrics

  • Query response times for cross-source queries
  • Storage growth rates
  • AI function execution times

Error Monitoring

  • Conversion failure rates
  • Database connection issues
  • AI function errors

Troubleshooting

Common Issues

  1. Conversion Failures

    • Check Terra webhook data format changes
    • Verify BasisEvent schema compatibility
    • Monitor device information completeness
  2. Query Performance

    • Check database indexing
    • Verify time range constraints
    • Monitor cross-source join complexity
  3. Data Inconsistencies

    • Validate timestamp synchronization
    • Check device deduplication logic
    • Verify health type mappings

Debug Tools

def debug_unified_data(user_id: str, start_date: str, end_date: str):
"""Debug tool for unified data issues."""
data_manager = get_unified_data_manager(user_id)

# Check data by source
apple_health_count = len(data_manager.query_cross_source_data(
health_types=[BasisHealthType.HEART_RATE],
start_time=datetime.strptime(start_date, '%Y-%m-%d'),
end_time=datetime.strptime(end_date, '%Y-%m-%d'),
device_sources=[MetaSource.APPLE_HEALTH]
))

terra_count = len(data_manager.query_cross_source_data(
health_types=[BasisHealthType.HEART_RATE],
start_time=datetime.strptime(start_date, '%Y-%m-%d'),
end_time=datetime.strptime(end_date, '%Y-%m-%d'),
device_sources=[MetaSource.TERRA]
))

print(f"Apple Health records: {apple_health_count}")
print(f"Terra records: {terra_count}")
print(f"Device registry: {len(data_manager.device_registry)}")

Conclusion

The unified data architecture enables seamless cross-source health data queries by:

  1. Normalizing all data sources to BasisSummaryV1 format
  2. Centralizing storage and querying through UnifiedDataManager
  3. Enhancing AI capabilities with cross-source tools
  4. Maintaining performance through optimized storage and indexing

This implementation addresses your original requirements:

  • ✅ Terra data flows and converts for unified querying
  • ✅ All data (health, notes, events) is queryable via Basis Copilot
  • ✅ Cross-device aggregation (Apple Watch + Oura Ring)
  • ✅ Event correlation (fasting timers with physiological data)
  • ✅ Seamless data flow from client to clinician view

The architecture is extensible for additional data sources and query types while maintaining backward compatibility with your existing systems.