diff --git a/docs/design/oura-import.md b/docs/design/oura-import.md index 3cf68a9bd..d5cf1980f 100644 --- a/docs/design/oura-import.md +++ b/docs/design/oura-import.md @@ -51,7 +51,7 @@ Endpoint names below are from model knowledge of the v2 API; **verify each again Other API facts to verify live at O2: OAuth2 endpoints (`cloud.ouraring.com/oauth/authorize`, `api.ouraring.com/oauth/token` ⚠), scopes (`daily heartrate workout tag session spo2 stress heart_health metabolic`), rate limit (historically 5000 requests / 5 min ⚠), the no-auth sandbox (`/v2/sandbox/usercollection/*` ⚠), personal-access-token deprecation status ⚠, webhook subscription API ⚠. -**Skeleton scope (implemented):** `daily_sleep`, `daily_readiness` (+ split-out `temperature_deviation` rows), `daily_resilience`, `daily_stress`, `daily_spo2`, `sleep`. That is exactly the "scores + stages" slice the day pages need and the AH mirror can't provide. +**Original skeleton scope:** `daily_sleep`, `daily_readiness` (+ split-out `temperature_deviation` rows), `daily_resilience`, `daily_stress`, `daily_spo2`, `sleep`. The shipped sync scope is broader; see §5 for the current OAuth scope set and §9a for the later granted-scope endpoint additions. ### Do we still need the pending Oura export? diff --git a/solstone/think/importers/apple_health.py b/solstone/think/importers/apple_health.py index 0f223870d..54a330968 100644 --- a/solstone/think/importers/apple_health.py +++ b/solstone/think/importers/apple_health.py @@ -245,7 +245,7 @@ class AppleHealthImporter: ) -> ImportResult: # Save mode fails closed before any parse or write, root-explicit # against the journal this call would actually write. - _gate_decision = enforce_pre_save_gate( + gate_decision = enforce_pre_save_gate( self, dry_run=dry_run, confirm_health_save=confirm_health_save, @@ -262,7 +262,7 @@ class AppleHealthImporter: summary=f"Dry run only: {preview.summary}", date_range=preview.date_range, ) - assert _gate_decision.raw_retention is not None + assert gate_decision.raw_retention is not None resolved_import_id = import_id or dt.datetime.now().strftime("%Y%m%d_%H%M%S") result = _save_export( @@ -272,7 +272,7 @@ class AppleHealthImporter: date_window=date_window, with_day_summaries=with_day_summaries, progress_callback=progress_callback, - retention=_gate_decision.raw_retention, + retention=gate_decision.raw_retention, ) return ImportResult( entries_written=result["entries_written"], diff --git a/solstone/think/importers/health_dedupe.py b/solstone/think/importers/health_dedupe.py index cdab45ce2..252145066 100644 --- a/solstone/think/importers/health_dedupe.py +++ b/solstone/think/importers/health_dedupe.py @@ -9,6 +9,7 @@ import datetime as dt import os import sqlite3 from collections.abc import Iterable +from contextlib import closing from dataclasses import dataclass from pathlib import Path from typing import Any @@ -120,9 +121,10 @@ def ensure_health_dedupe_db(journal_root: Path) -> Path: db_path = health_dedupe_db_path(journal_root) _precreate_private_sqlite_db(db_path) - with sqlite3.connect(db_path) as conn: - _ensure_health_dedupe_schema(conn) - _repair_sqlite_file_modes(db_path) + with closing(sqlite3.connect(db_path)) as conn: + with conn: + _ensure_health_dedupe_schema(conn) + _repair_sqlite_file_modes(db_path) return db_path @@ -137,7 +139,7 @@ def get_health_dedupe_record( return None _repair_sqlite_file_modes(db_path) - with sqlite3.connect(db_path) as conn: + with closing(sqlite3.connect(db_path)) as conn: conn.row_factory = sqlite3.Row row = conn.execute( "SELECT * FROM health_dedupe WHERE dedupe_key = ?", @@ -160,73 +162,74 @@ def upsert_health_dedupe_record( first_import_id = record.first_import_id or record.last_seen_import_id last_seen_import_id = record.last_seen_import_id or record.first_import_id - with sqlite3.connect(db_path) as conn: - _repair_sqlite_file_modes(db_path) - insert_result = conn.execute( - """ - INSERT INTO health_dedupe ( - dedupe_key, - source_family, - source_record_id, - record_type, - start_time, - end_time, - value_hash, - first_import_id, - last_seen_import_id, - normalized_ref, - raw_ref, - created_at, - updated_at + with closing(sqlite3.connect(db_path)) as conn: + with conn: + _repair_sqlite_file_modes(db_path) + insert_result = conn.execute( + """ + INSERT INTO health_dedupe ( + dedupe_key, + source_family, + source_record_id, + record_type, + start_time, + end_time, + value_hash, + first_import_id, + last_seen_import_id, + normalized_ref, + raw_ref, + created_at, + updated_at + ) + VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) + ON CONFLICT(dedupe_key) DO NOTHING + """, + ( + record.dedupe_key, + record.source_family, + record.source_record_id, + record.record_type, + record.start_time, + record.end_time, + record.value_hash, + first_import_id, + last_seen_import_id, + record.normalized_ref, + record.raw_ref, + now, + now, + ), + ) + if insert_result.rowcount == 1: + return True + + conn.execute( + """ + UPDATE health_dedupe + SET + source_record_id = COALESCE(?, source_record_id), + start_time = ?, + end_time = ?, + last_seen_import_id = ?, + value_hash = COALESCE(?, value_hash), + normalized_ref = COALESCE(?, normalized_ref), + raw_ref = COALESCE(?, raw_ref), + updated_at = ? + WHERE dedupe_key = ? + """, + ( + record.source_record_id, + record.start_time, + record.end_time, + last_seen_import_id, + record.value_hash, + record.normalized_ref, + record.raw_ref, + now, + record.dedupe_key, + ), ) - VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) - ON CONFLICT(dedupe_key) DO NOTHING - """, - ( - record.dedupe_key, - record.source_family, - record.source_record_id, - record.record_type, - record.start_time, - record.end_time, - record.value_hash, - first_import_id, - last_seen_import_id, - record.normalized_ref, - record.raw_ref, - now, - now, - ), - ) - if insert_result.rowcount == 1: - return True - - conn.execute( - """ - UPDATE health_dedupe - SET - source_record_id = COALESCE(?, source_record_id), - start_time = ?, - end_time = ?, - last_seen_import_id = ?, - value_hash = COALESCE(?, value_hash), - normalized_ref = COALESCE(?, normalized_ref), - raw_ref = COALESCE(?, raw_ref), - updated_at = ? - WHERE dedupe_key = ? - """, - ( - record.source_record_id, - record.start_time, - record.end_time, - last_seen_import_id, - record.value_hash, - record.normalized_ref, - record.raw_ref, - now, - record.dedupe_key, - ), - ) return False @@ -246,7 +249,7 @@ def upsert_health_dedupe_records( batch = tuple(records) now = dt.datetime.now(dt.UTC).isoformat().replace("+00:00", "Z") - with sqlite3.connect(db_path) as conn: + with closing(sqlite3.connect(db_path)) as conn: conn.execute("PRAGMA journal_mode=WAL") _repair_sqlite_file_modes(db_path) conn.execute("BEGIN") diff --git a/solstone/think/importers/oura.py b/solstone/think/importers/oura.py index 161172757..3b74afd74 100644 --- a/solstone/think/importers/oura.py +++ b/solstone/think/importers/oura.py @@ -109,7 +109,6 @@ from solstone.think.importers.sync import load_sync_state, save_sync_state logger = logging.getLogger(__name__) NORMALIZED_SCHEMA: Final = "solstone.health.oura.v1" -IMPORT_STREAM: Final = "import.oura" SYNC_BACKEND_NAME: Final = "oura" SYNC_STATE_SCHEMA: Final = "solstone.import_sync.oura.v1" # Journal-config section holding the Oura OAuth material (client_id, diff --git a/tests/test_health_dedupe.py b/tests/test_health_dedupe.py index cef0a8f49..b890c0a00 100644 --- a/tests/test_health_dedupe.py +++ b/tests/test_health_dedupe.py @@ -1,11 +1,13 @@ # SPDX-License-Identifier: AGPL-3.0-only # Copyright (c) 2026 sol pbc +import gc import os import sqlite3 import stat import time -from contextlib import contextmanager +import warnings +from contextlib import closing, contextmanager from pathlib import Path import pytest @@ -205,7 +207,7 @@ def test_upsert_health_dedupe_records_batches_in_wal_mode(tmp_path: Path): assert existing_row["normalized_ref"] == "import.apple_health/20260102/123000" assert new_row is not None assert new_row["source_record_id"] == "heart-rate-1" - with sqlite3.connect(health_dedupe_db_path(tmp_path)) as conn: + with closing(sqlite3.connect(health_dedupe_db_path(tmp_path))) as conn: journal_mode = conn.execute("PRAGMA journal_mode").fetchone()[0] assert journal_mode == "wal" @@ -217,7 +219,7 @@ def test_dedupe_db_and_wal_sidecars_are_private_while_connection_is_live( upsert_health_dedupe_records(tmp_path, _synthetic_dedupe_records(1)) db_path = health_dedupe_db_path(tmp_path) - with sqlite3.connect(db_path) as conn: + with closing(sqlite3.connect(db_path)) as conn: conn.execute("PRAGMA journal_mode=WAL") conn.execute( """ @@ -302,6 +304,23 @@ def test_batch_upsert_with_null_raw_ref_preserves_historic_raw_ref(tmp_path: Pat assert row["last_seen_import_id"] == "20260104_120000" +def test_dedupe_connections_close_without_resource_warnings(tmp_path: Path): + gc.collect() + with warnings.catch_warnings(record=True) as caught: + warnings.simplefilter("always", ResourceWarning) + + ensure_health_dedupe_db(tmp_path) + upsert_health_dedupe_record(tmp_path, _synthetic_dedupe_records(1)[0]) + upsert_health_dedupe_records(tmp_path, _synthetic_dedupe_records(2)) + assert get_health_dedupe_record(tmp_path, "sha256:synthetic-00000") is not None + gc.collect() + + resource_warnings = [ + warning for warning in caught if issubclass(warning.category, ResourceWarning) + ] + assert resource_warnings == [] + + def test_upsert_health_dedupe_records_handles_duplicate_keys_in_batch(tmp_path: Path): result = upsert_health_dedupe_records( tmp_path,