diff --git a/scripts/health/health_dashboard_action_worker.py b/scripts/health/health_dashboard_action_worker.py index 5efcc3a..3c40df3 100644 --- a/scripts/health/health_dashboard_action_worker.py +++ b/scripts/health/health_dashboard_action_worker.py @@ -1,175 +1,175 @@ #!/usr/bin/env python3 """Consume locally queued Dashboard actions outside the network-facing service.""" from __future__ import annotations import json import os import re import stat import subprocess import sys import sqlite3 import hashlib import hmac import secrets import shutil import tempfile from datetime import date, datetime from pathlib import Path from typing import Any from zoneinfo import ZoneInfo try: from dashboard_v5.patient_action_schema import REQUIRED_PATIENT_ACTION_COLUMNS from dashboard_v5.sprint6f_b_schema import assert_supplement_schema from dashboard_v5.supplement_contract import validate_supplement_payload from dashboard_v5.observation_contract import ( TRANSITIONS, validate_action as validate_observation_action, ) from dashboard_v5.sprint6f_c_schema import ( METHOD_VERSION, assert_schema as assert_observation_schema, ) from dashboard_v5.capture_contract import validate_capture_payload from dashboard_v5.medication_schema import assert_schema as assert_medication_schema from dashboard_v5.medication_contract import ( - ACTION_CONTRACT_VERSION as MEDICATION_ACTION_CONTRACT, + ACTION_CONTRACT_VERSIONS as MEDICATION_ACTION_CONTRACTS, resolve_action_preview, ) from dashboard_v5.capture_media import cleanup_expired, promote_attachment from dashboard_v5.media_validation import MediaInfo, create_safe_derivatives, media_runtime_self_test from dashboard_v5.sprint6i_b_schema import ( apply_schema as apply_media_schema, assert_schema as assert_media_schema, ) from dashboard_v5.sprint6h_b_schema import ( apply_schema as apply_capture_schema, assert_schema as assert_capture_schema, ) from dashboard_v5.document_review import ( CANDIDATE_ID_RE, ID_RE, TOKEN_RE, candidate_rows, extract_pages, load_quarantine, open_quarantine_source, normalize_search_text, section_similarity, validate_metadata, ) from dashboard_v5.document_originals import probe_original from dashboard_v5.sprint6i_a_schema import ( apply_schema as apply_document_schema, assert_schema as assert_document_schema, ) from dashboard_v5.sprint6i_c_schema import assert_schema as assert_reconciliation_schema from dashboard_v5.document_reconciliation import canonical_lab_pair, normalize_parameter, normalize_unit, normalize_value, reconcile_documents, transfer_preview from dashboard_v5.lab_review import QUALITATIVE, candidate_catalog_match, explicit_candidate_date, preview_revision_digest from dashboard_v5.metric_catalog_v2 import LAB_CATALOG_SPECS from dashboard_v5.nutrition_mapping_review import ( aggressive_identity as _normalize_food_identity, assignment_evidence_token, exact_identity as _exact_food_identity, resolve_catalog_assignment, resolve_group as resolve_mapping_group, ) except ModuleNotFoundError: # direct importlib fixture execution sys.path.insert(0, str(Path(__file__).resolve().parent)) from dashboard_v5.patient_action_schema import REQUIRED_PATIENT_ACTION_COLUMNS from dashboard_v5.sprint6f_b_schema import assert_supplement_schema from dashboard_v5.supplement_contract import validate_supplement_payload from dashboard_v5.observation_contract import ( TRANSITIONS, validate_action as validate_observation_action, ) from dashboard_v5.sprint6f_c_schema import ( METHOD_VERSION, assert_schema as assert_observation_schema, ) from dashboard_v5.capture_contract import validate_capture_payload from dashboard_v5.medication_schema import assert_schema as assert_medication_schema from dashboard_v5.medication_contract import ( - ACTION_CONTRACT_VERSION as MEDICATION_ACTION_CONTRACT, + ACTION_CONTRACT_VERSIONS as MEDICATION_ACTION_CONTRACTS, resolve_action_preview, ) from dashboard_v5.capture_media import cleanup_expired, promote_attachment from dashboard_v5.media_validation import MediaInfo, create_safe_derivatives, media_runtime_self_test from dashboard_v5.sprint6i_b_schema import ( apply_schema as apply_media_schema, assert_schema as assert_media_schema, ) from dashboard_v5.sprint6h_b_schema import ( apply_schema as apply_capture_schema, assert_schema as assert_capture_schema, ) from dashboard_v5.document_review import ( CANDIDATE_ID_RE, ID_RE, TOKEN_RE, candidate_rows, extract_pages, load_quarantine, open_quarantine_source, normalize_search_text, section_similarity, validate_metadata, ) from dashboard_v5.document_originals import probe_original from dashboard_v5.sprint6i_a_schema import ( apply_schema as apply_document_schema, assert_schema as assert_document_schema, ) from dashboard_v5.sprint6i_c_schema import assert_schema as assert_reconciliation_schema from dashboard_v5.document_reconciliation import canonical_lab_pair, normalize_parameter, normalize_unit, normalize_value, reconcile_documents, transfer_preview from dashboard_v5.lab_review import QUALITATIVE, candidate_catalog_match, explicit_candidate_date, preview_revision_digest from dashboard_v5.metric_catalog_v2 import LAB_CATALOG_SPECS from dashboard_v5.nutrition_mapping_review import ( aggressive_identity as _normalize_food_identity, assignment_evidence_token, exact_identity as _exact_food_identity, resolve_catalog_assignment, resolve_group as resolve_mapping_group, ) BASE = Path.home() / ".hermes" / "assets" / "Gesundheit" DEFAULT_DASHBOARD_DB = (BASE / "health_data.db").resolve() DASHBOARD_DB = ( Path(os.environ["HEALTH_DASHBOARD_DB"]).resolve() if os.environ.get("HEALTH_DASHBOARD_DB") else None ) ACTION_INBOX = Path( os.environ.get( "HEALTH_DASHBOARD_ACTION_INBOX", str(BASE / "runtime" / "dashboard-actions"), ) ).resolve() CAPTURE_MEDIA = Path( os.environ.get( "HEALTH_DASHBOARD_CAPTURE_MEDIA", str(BASE / "private-media" / "capture") ) ).resolve() CAPTURE_QUARANTINE = Path( os.environ.get( "HEALTH_DASHBOARD_CAPTURE_QUARANTINE", str(BASE / "runtime" / "capture-quarantine"), ) ).resolve() DOCUMENT_QUARANTINE = Path( os.environ.get( "HEALTH_DASHBOARD_DOCUMENT_QUARANTINE", str(BASE / "runtime" / "document-quarantine"), ) ).resolve() DOCUMENT_STORAGE = Path( os.environ.get( "HEALTH_DASHBOARD_DOCUMENT_STORAGE", str(BASE / "private-media" / "documents"), ) ).resolve() DASHBOARD_FILE = Path( @@ -966,378 +966,407 @@ def apply_supplement_action(payload: dict[str, Any]) -> None: """INSERT INTO supplement_intakes (plan_id,product,brand_variant,nutrient_key,amount,unit,status, occurred_at,composition_source,assignment_reliability,notes) VALUES(?,?,?,?,?,?,?,?,?,?,?)""", ( payload["plan_id"], payload["product"], payload["brand_variant"] or None, payload["nutrient_key"], payload["amount"], payload["unit"], payload["status"], payload["occurred_at"], payload["composition_source"], payload["assignment_reliability"], payload["notes"] or None, ), ) connection.commit() except Exception: connection.rollback() raise finally: connection.close() def apply_capture_action(payload: dict[str, Any]) -> str: if DASHBOARD_DB is None: raise RuntimeError("HEALTH_DASHBOARD_DB is required") payload = validate_capture_payload(payload) action_hash = hashlib.sha256(_canonical_json(payload).encode()).hexdigest() connection = sqlite3.connect(DASHBOARD_DB) connection.row_factory = sqlite3.Row promoted: list[dict[str, Any]] = [] try: connection.execute("PRAGMA foreign_keys=ON") apply_capture_schema(connection) assert_capture_schema(connection) apply_media_schema(connection) assert_media_schema(connection) connection.commit() connection.execute("BEGIN IMMEDIATE") existing = connection.execute( "SELECT action_hash,entry_id,request_version FROM capture_action_log WHERE idempotency_key=?", (payload["idempotency_key"],), ).fetchone() if existing: if ( existing["action_hash"] != action_hash or int(existing["request_version"]) != payload["request_version"] ): raise RuntimeError("conflicting idempotency key") connection.rollback() return str(existing["entry_id"]) now = datetime.now(LOCAL_TIMEZONE).isoformat(timespec="seconds") target_id = payload["corrects_entry_id"] or payload["withdraws_entry_id"] if target_id: target = connection.execute( "SELECT id,root_id,version,status,capture_type FROM capture_entries WHERE id=?", (target_id,), ).fetchone() if ( target is None or target["status"] != "active" or target["capture_type"] != payload["capture_type"] ): raise RuntimeError("version target is not active") root_id, version = str(target["root_id"]), int(target["version"]) + 1 connection.execute( "UPDATE capture_entries SET status=? WHERE id=?", ( "corrected" if payload["corrects_entry_id"] else "withdrawn", target_id, ), ) else: root_id, version = "", 1 entry_id = _opaque("cap_") if not root_id: root_id = entry_id data = payload["data"] medication_context: dict[str, Any] | None = None known_tables = { str(row[0]) for row in connection.execute("SELECT name FROM sqlite_master WHERE type='table'") } - if payload["capture_type"] == "medication" and data.get("contract") == MEDICATION_ACTION_CONTRACT: + if payload["capture_type"] == "medication" and data.get("contract") in MEDICATION_ACTION_CONTRACTS: assert_medication_schema(connection) expected_revision, prescription, planned, correction_target = resolve_action_preview( connection, payload ) if not hmac.compare_digest(expected_revision, data["preview_revision"]): raise RuntimeError("stale medication preview revision") medication_id = int(prescription["id"]) - duplicate = connection.execute( - """SELECT 1 FROM medication_administrations - WHERE business_revision IS NOT NULL AND medication_id=? - AND COALESCE(planned_event_id,-1)=COALESCE(?,-1) - AND COALESCE(occurred_at,'')=? - AND lower(trim(COALESCE(event_type,'')))=? - AND COALESCE(actual_dose_value,'')=? - AND COALESCE(actual_dose_unit,'')=? LIMIT 1""", - ( - medication_id, - int(planned["id"]) if planned is not None else None, - payload["occurred_at"], - data["status"], - data["actual_dose_value"], - data["actual_dose_unit"], - ), - ).fetchone() + if data["contract"] == "health.medication_action.v2": + duplicate = connection.execute( + """SELECT 1 FROM medication_administrations + WHERE business_revision IS NOT NULL AND medication_id=? + AND COALESCE(planned_event_id,-1)=COALESCE(?,-1) + AND COALESCE(occurred_at,'')=? + AND lower(trim(COALESCE(event_type,'')))=? + AND COALESCE(actual_quantity_value,'')=? + AND COALESCE(actual_dosage_form,'')=? + AND COALESCE(actual_strength,'')=? LIMIT 1""", + ( + medication_id, + int(planned["id"]) if planned is not None else None, + payload["occurred_at"], data["status"], + data["quantity_value"], data["dosage_form"], data["strength"], + ), + ).fetchone() + else: + duplicate = connection.execute( + """SELECT 1 FROM medication_administrations + WHERE business_revision IS NOT NULL AND medication_id=? + AND COALESCE(planned_event_id,-1)=COALESCE(?,-1) + AND COALESCE(occurred_at,'')=? + AND lower(trim(COALESCE(event_type,'')))=? + AND COALESCE(actual_dose_value,'')=? + AND COALESCE(actual_dose_unit,'')=? LIMIT 1""", + ( + medication_id, + int(planned["id"]) if planned is not None else None, + payload["occurred_at"], data["status"], + data["actual_dose_value"], data["actual_dose_unit"], + ), + ).fetchone() if duplicate and not data["duplicate_confirmed"]: raise RuntimeError("possible duplicate medication event") medication_context = { "medication_id": medication_id, "planned_event_id": int(planned["id"]) if planned is not None else None, "corrects_event_id": int(correction_target["id"]) if correction_target is not None else None, } elif payload["capture_type"] == "medication" and "medication_administrations" in known_tables: known_medications = { str(row[0]) for row in connection.execute( "SELECT DISTINCT medication_name FROM medication_administrations WHERE trim(COALESCE(medication_name,''))<>''" ) } if data["name"] not in known_medications: raise RuntimeError("medication is not in the known exact catalog") if payload["capture_type"] == "supplement" and "supplement_plans" in known_tables: known_supplements = { str(row[0]) for row in connection.execute( "SELECT DISTINCT product FROM supplement_plans WHERE trim(COALESCE(product,''))<>''" ) } if data["name"] not in known_supplements: raise RuntimeError("supplement is not in the known exact plan") - if payload["capture_type"] in {"medication", "supplement"} and data["status"] == "administered": + if ( + payload["capture_type"] in {"medication", "supplement"} + and data["status"] == "administered" + and data.get("contract") != "health.medication_action.v2" + ): if not (data["plan_value_confirmed"] or data["deviation_confirmed"]): raise RuntimeError("planned administration requires conscious confirmation") if payload["capture_type"] == "measurement": soft = { "systolic": (70, 220), "diastolic": (40, 140), "pulse": (35, 200), "weight": (30, 300), "temperature": (34, 42), "oxygen_saturation": (70, 100), "blood_glucose": (2, 30), } unusual = any( not soft[key][0] <= float(value) <= soft[key][1] for key, value in data["values"].items() ) if unusual and not data["plausibility_confirmed"]: raise RuntimeError("plausible outlier requires confirmation") duplicate = connection.execute( "SELECT 1 FROM capture_entries WHERE capture_type='measurement' AND occurred_at=? AND payload_json=? AND status='active'", (payload["occurred_at"], _canonical_json(data)), ).fetchone() if duplicate and not data["duplicate_confirmed"]: raise RuntimeError("duplicate measurement requires confirmation") for token in payload["attachments"]: attachment = promote_attachment(CAPTURE_QUARANTINE, CAPTURE_MEDIA, token) duplicate_media = connection.execute( """SELECT media_name,thumbnail_name,proxy_name,mime_type,media_kind,byte_size,width,height, frame_count,auxiliary_count,duration_seconds,container,codec,preview_sha256,proxy_sha256, preview_status,decoder_name,decoder_version,probe_version,unconfirmed_captured_at FROM capture_attachments WHERE sha256=? ORDER BY created_at LIMIT 1""", (attachment["sha256"],), ).fetchone() if duplicate_media: for duplicate_name in (attachment["media_name"], attachment["thumbnail_name"], attachment.get("proxy_name")): if duplicate_name: (CAPTURE_MEDIA / duplicate_name).unlink(missing_ok=True) attachment.update(dict(duplicate_media)) attachment["reused"] = True else: attachment["reused"] = False promoted.append(attachment) if sum(item["media_kind"] == "photo" for item in promoted) > 4 or sum(item["media_kind"] == "video" for item in promoted) > 1: raise RuntimeError("capture media count exceeded") entry_status = "withdrawn" if payload["withdraws_entry_id"] else "active" connection.execute( "INSERT INTO capture_entries(id,root_id,version,capture_type,occurred_at,ended_at,payload_json,status,corrects_entry_id,withdraws_entry_id,created_at) VALUES(?,?,?,?,?,?,?,?,?,?,?)", ( entry_id, root_id, version, payload["capture_type"], payload["occurred_at"], payload["ended_at"], _canonical_json(data), entry_status, payload["corrects_entry_id"], payload["withdraws_entry_id"], now, ), ) for attachment in promoted: connection.execute( """INSERT INTO capture_attachments( id,entry_id,sha256,mime_type,media_kind,byte_size,width,height,frame_count,auxiliary_count, duration_seconds,container,codec,media_name,thumbnail_name,proxy_name,preview_sha256,proxy_sha256, preview_status,decoder_name,decoder_version,probe_version,unconfirmed_captured_at,review_status, description,body_region,created_at) VALUES(?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)""", ( attachment["id"], entry_id, attachment["sha256"], attachment["mime_type"], attachment["media_kind"], attachment["byte_size"], attachment["width"], attachment["height"], attachment["frame_count"], attachment["auxiliary_count"], attachment.get("duration_seconds"), attachment.get("container"), attachment.get("codec"), attachment["media_name"], attachment["thumbnail_name"], attachment.get("proxy_name"), attachment.get("preview_sha256"), attachment.get("proxy_sha256"), attachment["preview_status"], attachment["decoder_name"], attachment["decoder_version"], attachment.get("probe_version"), attachment.get("unconfirmed_captured_at"), "unverified", data.get("description") or None, data.get("body_region") or None, now, ), ) if payload["corrects_entry_id"] and not promoted: columns = "sha256,mime_type,media_kind,byte_size,width,height,frame_count,auxiliary_count,duration_seconds,container,codec,media_name,thumbnail_name,proxy_name,preview_sha256,proxy_sha256,preview_status,decoder_name,decoder_version,probe_version,unconfirmed_captured_at,review_status,explicit_pair_id,description,body_region" for old in connection.execute(f"SELECT {columns} FROM capture_attachments WHERE entry_id=?", (payload["corrects_entry_id"],)).fetchall(): prefix = "vid_" if old[2] == "video" else "img_" connection.execute( f"INSERT INTO capture_attachments(id,entry_id,{columns},created_at) VALUES(?,?{',?' * len(old)},?)", (_opaque(prefix), entry_id, *old, now), ) day = payload["occurred_at"].split("T", 1)[0] available_tables = { str(row[0]) for row in connection.execute("SELECT name FROM sqlite_master WHERE type='table'") } if not payload["withdraws_entry_id"]: if payload["capture_type"] == "symptom" and "symptom_log" in available_tables: severity = None if data["intensity"] == "unknown" else data["intensity"] label = data["title"] or data["category"] context = _canonical_json( { "source": "mobile_capture", "category": data["category"], "count": data["count"], "body_region": data["body_region"], "ongoing": data["ongoing"], "entry_id": entry_id, } ) symptom_columns = { str(row[1]) for row in connection.execute("PRAGMA table_info(symptom_log)") } if {"occurred_at", "onset_at", "duration_minutes"}.issubset(symptom_columns): connection.execute( "INSERT INTO symptom_log(datum,symptom,schwergrad,kontext,notizen,occurred_at,onset_at,duration_minutes) VALUES(?,?,?,?,?,?,?,NULL)", ( day, label, severity, context, data["note"] or None, payload["occurred_at"], payload["occurred_at"], ), ) else: connection.execute( "INSERT INTO symptom_log(datum,symptom,schwergrad,kontext,notizen) VALUES(?,?,?,?,?)", (day, label, severity, context, data["note"] or None), ) elif payload["capture_type"] == "medication" and "medication_administrations" in available_tables: if medication_context is not None: - actual_dose = " ".join( - part for part in (data["actual_dose_value"], data["actual_dose_unit"]) if part - ) - legacy_route = data["route_original"] or data["route_normalized"] - connection.execute( - """INSERT INTO medication_administrations( - datum,medication_name,dose,route,event_type,scheduled_next_date,notes,source,occurred_at, - medication_id,planned_event_id,planned_dose_value,planned_dose_unit,actual_dose_value, - actual_dose_unit,route_original,route_normalized,injection_region,injection_side, - injection_detail,lot_number,corrects_event_id,corrected_target_status,correction_reason, - business_revision) VALUES(?,?,?,?,?,NULL,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)""", - ( - day, - data["name"], - actual_dose or None, - legacy_route or None, - data["status"], - data["note"] or None, - "dashboard_v5_medication_action", - payload["occurred_at"], - medication_context["medication_id"], - medication_context["planned_event_id"], - data["planned_dose_value"] or None, - data["planned_dose_unit"] or None, - data["actual_dose_value"] or None, - data["actual_dose_unit"] or None, - data["route_original"] or None, - data["route_normalized"] or None, - data["injection_region"] or None, - data["injection_side"] or None, - data["injection_detail"] or None, - data["lot_number"] or None, - medication_context["corrects_event_id"], - data["corrected_target_status"] or None, - data["correction_reason"] or None, - action_hash, - ), - ) + if data["contract"] == "health.medication_action.v2": + legacy_dose = f"{data['quantity_value']} {data['dosage_form']} · {data['strength']}" + legacy_route = data["route_original"] or data["route_normalized"] + connection.execute( + """INSERT INTO medication_administrations( + datum,medication_name,dose,route,event_type,scheduled_next_date,notes,source,occurred_at, + medication_id,planned_event_id,route_original,route_normalized,injection_region, + injection_side,injection_detail,business_revision,actual_quantity_value, + actual_dosage_form,actual_strength) + VALUES(?,?,?,?,?,NULL,?,?,?,?,?,?,?,?,?,?,?,?,?,?)""", + ( + day, data["name"], legacy_dose, legacy_route or None, + data["status"], data["note"] or None, + "dashboard_v5_medication_action_v2", payload["occurred_at"], + medication_context["medication_id"], medication_context["planned_event_id"], + data["route_original"] or None, data["route_normalized"] or None, + data["injection_region"] or None, data["injection_side"] or None, + data["injection_detail"] or None, action_hash, + data["quantity_value"], data["dosage_form"], data["strength"], + ), + ) + else: + actual_dose = " ".join( + part for part in (data["actual_dose_value"], data["actual_dose_unit"]) if part + ) + legacy_route = data["route_original"] or data["route_normalized"] + connection.execute( + """INSERT INTO medication_administrations( + datum,medication_name,dose,route,event_type,scheduled_next_date,notes,source,occurred_at, + medication_id,planned_event_id,planned_dose_value,planned_dose_unit,actual_dose_value, + actual_dose_unit,route_original,route_normalized,injection_region,injection_side, + injection_detail,lot_number,corrects_event_id,corrected_target_status,correction_reason, + business_revision) VALUES(?,?,?,?,?,NULL,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)""", + ( + day, data["name"], actual_dose or None, legacy_route or None, + data["status"], data["note"] or None, + "dashboard_v5_medication_action", payload["occurred_at"], + medication_context["medication_id"], medication_context["planned_event_id"], + data["planned_dose_value"] or None, data["planned_dose_unit"] or None, + data["actual_dose_value"] or None, data["actual_dose_unit"] or None, + data["route_original"] or None, data["route_normalized"] or None, + data["injection_region"] or None, data["injection_side"] or None, + data["injection_detail"] or None, data["lot_number"] or None, + medication_context["corrects_event_id"], data["corrected_target_status"] or None, + data["correction_reason"] or None, action_hash, + ), + ) else: dose = " ".join( str(part) for part in (data["amount"], data["unit"]) if part not in {None, ""} ) medication_columns = {str(row[1]) for row in connection.execute("PRAGMA table_info(medication_administrations)")} if "occurred_at" in medication_columns: connection.execute( "INSERT INTO medication_administrations(datum,medication_name,dose,route,event_type,notes,source,occurred_at) VALUES(?,?,?,?,?,?,?,?)", (day, data["name"], dose or None, data["route"] or None, data["status"], data["note"] or None, "dashboard_v5_mobile_capture", payload["occurred_at"]), ) else: connection.execute( "INSERT INTO medication_administrations(datum,medication_name,dose,route,event_type,notes,source) VALUES(?,?,?,?,?,?,?)", (day, data["name"], dose or None, data["route"] or None, data["status"], data["note"] or None, "dashboard_v5_mobile_capture"), ) elif payload["capture_type"] == "supplement" and "supplement_intakes" in available_tables: connection.execute( "INSERT INTO supplement_intakes(plan_id,product,brand_variant,nutrient_key,amount,unit,status,occurred_at,composition_source,assignment_reliability,notes) VALUES(NULL,?,NULL,'other',?,?,?,?, 'user_documented','documented',?)", ( data["name"], data["amount"], data["unit"], data["status"], payload["occurred_at"], data["note"] or None, ), ) elif payload["capture_type"] == "event" and "health_events" in available_tables: event_columns = { str(row[1]) for row in connection.execute("PRAGMA table_info(health_events)") } event_kind = data.get("event_kind") if event_kind == "sauna": category = "sauna_recovery" parameter = "Sauna" value = data["duration_minutes"] unit = "min" source = "dashboard_v5_manual_capture" intensity = None elif event_kind == "training": category = "manual_training" parameter = data["activity_type"] value = data["duration_minutes"] unit = "min" source = "dashboard_v5_manual_capture" intensity = None else: category = data["category"] parameter = data["title"] intensity = ( None if data["intensity"] == "unknown" else data["intensity"] ) value = intensity unit = None source = "dashboard_v5_mobile_capture" if {"occurred_at", "intensity"}.issubset(event_columns): connection.execute( "INSERT INTO health_events(date,category,parameter,value,unit,source,notes,occurred_at,intensity) VALUES(?,?,?,?,?,?,?,?,?)", ( day, category, parameter, value, unit, source, data["note"] or None, payload["occurred_at"], intensity, ), ) else: connection.execute( "INSERT INTO health_events(date,category,parameter,value,unit,source,notes) VALUES(?,?,?,?,?,?,?)", ( day, category, parameter, diff --git a/scripts/health/health_dashboard_server.py b/scripts/health/health_dashboard_server.py index 58e878b..6fefe07 100644 --- a/scripts/health/health_dashboard_server.py +++ b/scripts/health/health_dashboard_server.py @@ -1,118 +1,118 @@ #!/usr/bin/env python3 """Serve the local Health Dashboard with strict routes and privacy headers.""" from __future__ import annotations import base64 import binascii import hashlib import json import mimetypes import os import re import secrets import sqlite3 import stat import threading import time from datetime import date, datetime, timedelta from http.cookies import CookieError, SimpleCookie from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer from pathlib import Path from typing import Any from urllib.parse import parse_qs, quote, unquote, urlparse from zoneinfo import ZoneInfo from dashboard_v5.read_api import ( APIError, connect_read_only, dispatch_api, _resolve_record_document, ) from dashboard_v5.document_originals import configured_original_roots, probe_original from dashboard_v5.supplement_contract import validate_supplement_payload from dashboard_v5.observation_contract import validate_action as validate_observation_action from dashboard_v5.capture_contract import validate_capture_payload from dashboard_v5.medication_schema import assert_schema as assert_medication_schema from dashboard_v5.medication_contract import ( - ACTION_CONTRACT_VERSION as MEDICATION_ACTION_CONTRACT, + ACTION_CONTRACT_VERSIONS as MEDICATION_ACTION_CONTRACTS, resolve_action_preview, ) from dashboard_v5.capture_media import MAX_ORIGINAL_BYTES as MAX_CAPTURE_MEDIA_BYTES, quarantine_bytes from dashboard_v5.media_validation import media_runtime_self_test from dashboard_v5.document_review import discard_quarantine, quarantine_upload, validate_metadata BASE = Path.home() / ".hermes" / "assets" / "Gesundheit" API_DB = ( Path(os.environ["HEALTH_DASHBOARD_DB"]).resolve() if os.environ.get("HEALTH_DASHBOARD_DB") else None ) DB = BASE / "health_data.db" REPORTS = BASE / "reports" DASHBOARD_FILE = Path(os.environ["HEALTH_DASHBOARD_FILE"]).resolve() DASHBOARD_V5_FILE = ( Path(os.environ["HEALTH_DASHBOARD_V5_FILE"]).resolve() if os.environ.get("HEALTH_DASHBOARD_V5_FILE") else None ) ASSET_DIR = Path( os.environ.get( "HEALTH_DASHBOARD_ASSET_DIR", str(BASE / "scripts" / "assets" / "health-assets"), ) ).resolve() BIND_HOST = os.environ.get("HEALTH_DASHBOARD_HOST", "127.0.0.1") BIND_PORT = int(os.environ.get("HEALTH_DASHBOARD_PORT", "8014")) EXTERNAL_SCHEME = os.environ.get("HEALTH_DASHBOARD_SCHEME", "http").strip().lower() if EXTERNAL_SCHEME not in {"http", "https"}: raise RuntimeError("HEALTH_DASHBOARD_SCHEME must be http or https") EXTERNAL_PORT = int(os.environ.get("HEALTH_DASHBOARD_PUBLIC_PORT", str(BIND_PORT))) if not 1 <= EXTERNAL_PORT <= 65535: raise RuntimeError("HEALTH_DASHBOARD_PUBLIC_PORT must be between 1 and 65535") ALLOWED_HOSTS = { value.strip().lower().rstrip(".") for value in os.environ.get( "HEALTH_DASHBOARD_ALLOWED_HOSTS", "127.0.0.1,localhost" ).split(",") if value.strip() } ROUTE = "/health-dashboard" V5_ROUTE = "/health-dashboard-v5" ASSET_ROUTES = { "/health-assets/chart.umd.min.js": ( "chart.umd.min.js", "text/javascript; charset=utf-8", ), "/health-assets/dashboard-v5.js": ( "dashboard-v5.js", "text/javascript; charset=utf-8", ), "/health-assets/dashboard-v5-nutrition.js": ( "dashboard-v5-nutrition.js", "text/javascript; charset=utf-8", ), "/health-assets/dashboard-v5-range.js": ( "dashboard-v5-range.js", "text/javascript; charset=utf-8", ), "/health-assets/dashboard-v5-global-search.js": ( "dashboard-v5-global-search.js", "text/javascript; charset=utf-8", ), "/health-assets/dashboard-v5-api-explorer.js": ( "dashboard-v5-api-explorer.js", "text/javascript; charset=utf-8", ), "/health-assets/dashboard-v5-associations.js": ( "dashboard-v5-associations.js", "text/javascript; charset=utf-8", ), "/health-assets/dashboard-v5-observations.js": ( "dashboard-v5-observations.js", "text/javascript; charset=utf-8", ), "/health-assets/dashboard-v5-echarts.js": ( "dashboard-v5-echarts.js", "text/javascript; charset=utf-8", @@ -1116,160 +1116,181 @@ def write_action_payload(payload: dict[str, object]) -> str: if receipt.is_file() and not receipt.is_symlink(): return token if sum(1 for _ in ACTION_INBOX.glob("*.json")) >= MAX_PENDING_ACTIONS: raise QueueFullError("private action queue is full") temporary = ACTION_INBOX / f".{token}.tmp" destination = ACTION_INBOX / f"{token}.json" descriptor = os.open( temporary, os.O_WRONLY | os.O_CREAT | os.O_EXCL | os.O_NOFOLLOW, 0o600, ) try: with os.fdopen(descriptor, "wb") as handle: handle.write(packed) handle.flush() os.fsync(handle.fileno()) temporary.replace(destination) directory = os.open(ACTION_INBOX, os.O_RDONLY | os.O_DIRECTORY) try: os.fsync(directory) finally: os.close(directory) except Exception: temporary.unlink(missing_ok=True) raise return token def validate_document_review_submission(form: dict[str, list[str]]) -> tuple[dict[str, object], str, str]: if set(form) != {"csrf_token", "payload", "return_to", "return_document"} or any(len(values) != 1 for values in form.values()) or form["return_to"][0] not in {"v5", "v5_labs"}: raise ValueError("invalid document review submission") return_document=form["return_document"][0] if not re.fullmatch(r"api-document-[a-f0-9]{24}",return_document):raise ValueError("invalid return document") payload = json.loads(form["payload"][0], object_pairs_hook=reject_duplicate_object_pairs) required = {"version","action","operation","document_id","target_id","value","unit","decision","metadata"} optional = {"expected_candidate_revision", "preview_revision"} if not isinstance(payload, dict) or not required.issubset(payload) or not set(payload).issubset(required | optional) or payload.get("version") != 1 or payload.get("action") != "document_review": raise ValueError("invalid document review payload") if not isinstance(payload.get("document_id"), str) or not re.fullmatch(r"doc_[a-f0-9]{24}", payload["document_id"]): raise ValueError("invalid document review id") if payload.get("operation") not in {"metadata_review","text_correction","candidate_decision","candidate_remap","candidate_transfer","content_review","discard_document","retry_extraction","original_review","page_review"}: raise ValueError("invalid document operation") if not isinstance(payload.get("target_id"), str) or len(payload["target_id"]) > 64: raise ValueError("invalid document target") if payload["operation"] in {"text_correction","page_review"} and not re.fullmatch(r"page_[1-9][0-9]{0,2}",payload["target_id"]): raise ValueError("invalid document page target") if payload["operation"] in {"candidate_decision","candidate_remap","candidate_transfer"} and not re.fullmatch(r"cand_[a-f0-9]{24}",payload["target_id"]): raise ValueError("invalid document candidate target") expected_revision = payload.get("expected_candidate_revision") preview_revision = payload.get("preview_revision") if form["return_to"][0] == "v5_labs" and payload["operation"] in {"candidate_decision","candidate_remap","candidate_transfer"} and (expected_revision is None or preview_revision is None): raise ValueError("candidate preview binding required") if form["return_to"][0] == "v5_labs" and payload["operation"] == "candidate_decision" and payload.get("decision") == "corrected_confirmed": raise ValueError("laboratory correction requires a regenerated preview") if expected_revision is not None and (type(expected_revision) is not int or not 1 <= expected_revision <= 2_147_483_647): raise ValueError("invalid candidate revision") if preview_revision is not None and (not isinstance(preview_revision, str) or not re.fullmatch(r"[a-f0-9]{64}", preview_revision)): raise ValueError("invalid preview revision") if not isinstance(payload.get("value"), str) or len(payload["value"]) > 4000 or not isinstance(payload.get("unit"), str) or len(payload["unit"]) > 40 or not isinstance(payload.get("decision"), str) or len(payload["decision"]) > 30: raise ValueError("invalid document review value") metadata = payload.get("metadata") if payload["operation"] == "metadata_review": payload["metadata"] = validate_metadata(metadata) elif not isinstance(metadata, dict): raise ValueError("invalid document metadata") return payload, form["return_to"][0], return_document def write_symptom_checkin(day: str, scores: dict[str, int], notes: str) -> str: return write_action_payload( { "version": 1, "action": "symptom_checkin", "date": day, "scores": scores, "notes": notes, } ) +def _medication_preview_error(error: Exception) -> tuple[int, dict[str, Any]]: + reason = str(error) + rules = ( + ("invalid medication action shape", 422, "payload_invalid", "medication", "Die Medikationsangaben sind unvollständig oder ungültig."), + ("invalid medication text", 422, "medication_text_invalid", "medication", "Name, Stärke oder Notiz enthält eine unzulässige Angabe."), + ("requires a real plan", 422, "planned_event_required", "capture_mode", "Für diesen Modus ist kein belastbarer Plantermin verknüpft."), + ("lacks a structured administration preset", 422, "planned_event_unstructured", "capture_mode", "Der Plantermin enthält keine revisionssichere Mengenangabe."), + ("differs from verified preset", 422, "preset_deviation_unconfirmed", "deviation_confirmed", "Die Eingabe weicht von der verifizierten Vorauswahl ab und muss bewusst bestätigt werden."), + ("administration preset changed", 409, "preset_changed", "preset", "Die verifizierte Vorauswahl hat sich geändert. Bitte Vorschau neu laden."), + ("stale medication preview", 409, "preview_changed", "preview", "Die Vorschau ist nicht mehr aktuell. Bitte erneut prüfen."), + ("already consumed", 409, "planned_event_consumed", "capture_mode", "Der Plantermin wurde bereits dokumentiert."), + ("explicitly active prescription", 422, "prescription_not_active_for_plan", "capture_mode", "Ein Plantermin kann nur gegen eine ausdrücklich aktive Verordnung dokumentiert werden."), + ) + for fragment, status_code, code, field, message in rules: + if fragment in reason: + return status_code, {"error": {"kind": "validation", "code": code, "field": field, "message": message}} + if isinstance(error, ValueError): + return 422, {"error": {"kind": "validation", "code": "medication_validation_failed", "field": "medication", "message": "Die Medikationsangaben sind ungültig. Bitte die Felder prüfen."}} + return 503, {"error": {"kind": "technical", "code": "medication_preview_unavailable", "field": "", "message": "Die Erfassung ist technisch nicht verfügbar. Es wurde nichts vorgemerkt."}} + + class Handler(BaseHTTPRequestHandler): def send_error( self, code: int, message: str | None = None, explain: str | None = None, ) -> None: request_path = urlparse(getattr(self, "path", "")).path if request_path.startswith("/api/v1/"): send_body = getattr(self, "command", "") != "HEAD" if not host_is_allowed(self.headers.get("Host", "")): self._send_api_error(421, "invalid_host", send_body) return status = 405 if code == 501 else code error_code = "read_only_endpoint" if code == 501 else "request_rejected" self._send_api_error(status, error_code, send_body) return super().send_error(code, message, explain) def _reject_unsupported(self) -> None: if not host_is_allowed(self.headers.get("Host", "")): if urlparse(self.path).path.startswith("/api/v1/"): self._send_api_error(421, "host_not_allowed", self.command != "HEAD") else: self.send_error(421) return if urlparse(self.path).path.startswith("/api/v1/"): self._send_api_error(405, "read_only_endpoint", self.command != "HEAD") return self.send_error(405) do_OPTIONS = _reject_unsupported # noqa: N815 do_PUT = _reject_unsupported # noqa: N815 do_DELETE = _reject_unsupported # noqa: N815 do_PATCH = _reject_unsupported # noqa: N815 do_TRACE = _reject_unsupported # noqa: N815 do_CONNECT = _reject_unsupported # noqa: N815 def do_HEAD(self) -> None: # noqa: N802 if not host_is_allowed(self.headers.get("Host", "")): if urlparse(self.path).path.startswith("/api/v1/"): self._send_api_error(421, "host_not_allowed", False) else: self.send_error(421) return self._handle(send_body=False) def do_GET(self) -> None: # noqa: N802 if not host_is_allowed(self.headers.get("Host", "")): if urlparse(self.path).path.startswith("/api/v1/"): self._send_api_error(421, "host_not_allowed", True) else: self.send_error(421) return self._handle(send_body=True) def do_POST(self) -> None: # noqa: N802 if not host_is_allowed(self.headers.get("Host", "")): if urlparse(self.path).path.startswith("/api/v1/"): self._send_api_error(421, "host_not_allowed", True) else: self.send_error(421) return api_path = urlparse(self.path).path if api_path.startswith("/api/v1/"): if api_path == "/api/v1/browser-session": self._handle_browser_session(True) elif api_path == CAPTURE_UPLOAD_ROUTE: self._handle_capture_upload() elif api_path == DOCUMENT_UPLOAD_ROUTE: self._handle_document_upload() else: self._send_api_error(405, "read_only_endpoint", True) return action_path = urlparse(self.path).path if action_path not in { CHECKIN_ROUTE, NUTRITION_MAPPING_ROUTE, SYMPTOM_EVENT_ROUTE, MEDICATION_EVENT_ROUTE, @@ -1279,201 +1300,203 @@ class Handler(BaseHTTPRequestHandler): DOCUMENT_REVIEW_ROUTE, CAPTURE_ROUTE, MEDICATION_PREVIEW_ROUTE, }: self.send_error(404) return if not browser_session_is_authenticated(self.headers.get("Cookie", "")): if "application/json" in self.headers.get("Accept", ""): self._send_api_error(401, "authentication_required", True) else: self.send_error(401) return origin = self.headers.get("Origin", "") host = self.headers.get("Host", "") if not origin_matches_request(origin, host, int(getattr(self.server, "server_port", EXTERNAL_PORT))): self.send_error(403) return if self.headers.get_content_type() != "application/x-www-form-urlencoded": self.send_error(415) return try: length = int(self.headers.get("Content-Length", "0")) except ValueError: self.send_error(400) return maximum_body = 8192 if action_path == DOCUMENT_REVIEW_ROUTE else 4096 if length < 1 or length > maximum_body: self.send_error(413) return try: form = parse_qs( self.rfile.read(length).decode("utf-8"), keep_blank_values=True, strict_parsing=True, max_num_fields=20, ) except (UnicodeError, ValueError): self.send_error(400) return token = (form.get("csrf_token") or [""])[0] cookie = SimpleCookie(self.headers.get("Cookie", "")) cookie_token = cookie.get( "health_capture_csrf" if action_path in {CAPTURE_ROUTE, MEDICATION_PREVIEW_ROUTE} else "health_csrf" ) if ( not token or not cookie_token or not secrets.compare_digest(token, cookie_token.value) ): self.send_error(403) return if action_path != MEDICATION_PREVIEW_ROUTE and not consume_csrf_token(token): self.send_error(403) return mapping_queue_key_value = "" medication_preview_revision = "" return_document = "" payload: dict[str, object] = {} try: if action_path == CHECKIN_ROUTE: day, scores, notes, return_to = validate_symptom_submission(form) write_symptom_checkin(day, scores, notes) elif action_path == NUTRITION_MAPPING_ROUTE: payload, return_to = validate_mapping_submission(form) mapping_queue_key_value = str(payload["queue_key"]) write_action_payload(payload) elif action_path == MEDICATION_PREVIEW_ROUTE: if ( set(form) != {"csrf_token", "payload", "return_to"} or any(len(values) != 1 for values in form.values()) or form["return_to"][0] != "v5" or API_DB is None ): raise ValueError("invalid medication preview submission") payload = validate_capture_payload(json.loads(form["payload"][0])) preview_data = payload.get("data") if ( payload["capture_type"] != "medication" or not isinstance(preview_data, dict) - or preview_data.get("contract") != MEDICATION_ACTION_CONTRACT + or preview_data.get("contract") not in MEDICATION_ACTION_CONTRACTS ): raise ValueError("invalid medication preview contract") preview_connection = connect_read_only(API_DB) try: assert_medication_schema(preview_connection) medication_preview_revision = resolve_action_preview( preview_connection, payload )[0] - except RuntimeError as error: - raise ValueError("medication preview context unavailable") from error finally: preview_connection.close() return_to = "v5" elif action_path == CAPTURE_ROUTE: if ( set(form) != {"csrf_token", "payload", "return_to"} or any(len(values) != 1 for values in form.values()) or form["return_to"][0] != "v5" ): raise ValueError("invalid capture submission") payload = validate_capture_payload(json.loads(form["payload"][0])) return_to = "v5" write_action_payload(payload) elif action_path == DOCUMENT_REVIEW_ROUTE: payload, return_to, return_document = validate_document_review_submission(form) write_action_payload(payload) elif action_path == OBSERVATION_ROUTE: payload, return_to = validate_observation_submission(form) write_action_payload(payload) else: payload, return_to = validate_patient_action(form, action_path) write_action_payload(payload) except IdempotencyConflictError: if "application/json" in self.headers.get("Accept", ""): self._send_api_error(409, "idempotency_conflict", True) else: self.send_error(409) return - except ValueError: - self.send_error(400) + except (ValueError, RuntimeError, sqlite3.Error) as error: + if action_path == MEDICATION_PREVIEW_ROUTE: + status_code, body = _medication_preview_error(error) + self._send_api_json(status_code, body, True) + else: + self.send_error(400) return except QueueFullError: self.send_error(503) return except OSError: self.send_error(500) return if action_path == MEDICATION_PREVIEW_ROUTE: self._send_api_json( 200, {"status": "preview", "preview_revision": medication_preview_revision}, True, ) return if action_path in {CHECKIN_ROUTE, NUTRITION_MAPPING_ROUTE} and "application/json" in self.headers.get( "Accept", "" ): self._send_api_json( 202, {"status": "queued"}, True, ) return if action_path == CAPTURE_ROUTE and "application/json" in self.headers.get( "Accept", "" ): self._send_api_json( 202, {"status": "queued", "idempotency_key": payload["idempotency_key"]}, True, ) return if return_to in {"v5", "v5_labs"}: receipt = quote(issue_queue_receipt(), safe="") if action_path == NUTRITION_MAPPING_ROUTE: queue_key = quote(mapping_queue_key_value, safe="") target = ( V5_ROUTE + "?view=nutrition&workspace=mapping&map_status=open" + f"&map_sort=frequency&mapping={queue_key}&queued={receipt}" ) elif action_path == DOCUMENT_REVIEW_ROUTE and return_to == "v5_labs": target = V5_ROUTE + f"?view=record&tab=labs&queued={receipt}#lab-review" elif action_path == DOCUMENT_REVIEW_ROUTE: target = V5_ROUTE + f"?view=record&tab=documents&document={quote(return_document,safe='')}&queued={receipt}#next-document-decision" elif action_path == OBSERVATION_ROUTE: target = ( V5_ROUTE + f"?view=explorer&explorer_area=questions&queued={receipt}#observation-status" ) else: target = V5_ROUTE + f"?queued={receipt}#capture-status" else: target = ROUTE + "?queued=1#symptome" self.send_response(303) self.send_header("Location", target) self.send_header("Cache-Control", "no-store") self.end_headers() def _v5_principal_authenticated(self) -> bool: if API_DB is None: return True try: token = load_api_token() basic_password = load_dashboard_basic_password() except (OSError, UnicodeError): return False return dashboard_principal_is_authenticated( self.headers.get("Authorization", ""), token, basic_password ) def _send_dashboard_auth_required(self) -> None: self.send_response(401) self.send_header("WWW-Authenticate", 'Basic realm="Health Dashboard V5"') self.send_header("Cache-Control", "no-store") self.send_header("Pragma", "no-cache") self.send_header("X-Content-Type-Options", "nosniff") self.send_header("Content-Length", "0") self.end_headers() __HERMES_CWD_8d46a20096ed__/home/agent/.hermes/repos/HealthManager__HERMES_CWD_8d46a20096ed__