From b1a6420a4e71ddce5b18e7a2c2cc25d50f36a823 Mon Sep 17 00:00:00 2001 From: Jer Miller Date: Mon, 13 Jul 2026 07:31:25 -0600 Subject: [PATCH] fix(importers): close health dedupe sqlite connections with sqlite3.connect(...) commits or rolls back but never closes, so every dedupe write leaked an open DB handle, visible as ResourceWarnings under CI, and repeated import or sync runs could accumulate DB/WAL handles. Connections are now wrapped in contextlib.closing with the transaction scope preserved: a nested with conn: where the implicit commit was relied on, and the explicit BEGIN/commit/rollback left intact on the batch path. Also removes a dead duplicate import.oura stream literal that sat outside the registry, renames the now load-bearing Apple Health gate decision off its placeholder underscore name, and fixes a stale Oura design-doc line that understated the implemented endpoint set. (cherry picked from commit b9299a646a6ad23d7990460f3c6b377fbbe81242) --- docs/design/oura-import.md | 2 +- solstone/think/importers/apple_health.py | 6 +- solstone/think/importers/health_dedupe.py | 145 +++++++++++----------- solstone/think/importers/oura.py | 1 - tests/test_health_dedupe.py | 25 +++- 5 files changed, 100 insertions(+), 79 deletions(-) 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, -- 2.51.2