#!/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,
        bind_verified_hyrimoz_preset,
    )
    from dashboard_v5.medication_contract import (
        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,
        bind_verified_hyrimoz_preset,
    )
    from dashboard_v5.medication_contract import (
        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(
    os.environ.get(
        "HEALTH_DASHBOARD_FILE",
        str(BASE / "reports" / "health_dashboard_v4.html"),
    )
).resolve()
DASHBOARD_V5_FILE = (
    Path(os.environ["HEALTH_DASHBOARD_V5_FILE"]).resolve()
    if os.environ.get("HEALTH_DASHBOARD_V5_FILE")
    else None
)
SCRIPT_DIR = Path(__file__).resolve().parent
QUICK_ADD = SCRIPT_DIR / "health_symptom_quick_add.py"
DASHBOARD_V4 = SCRIPT_DIR / "health_dashboard_v4.py"
DASHBOARD_V5 = SCRIPT_DIR / "health_dashboard_v5.py"
FIELDS = ("aphthen", "gi", "fatigue", "skin", "eyes", "joints", "vascular")
NAME_RE = re.compile(r"[0-9a-f]{32}\.json")
LOCAL_TIMEZONE = ZoneInfo("Europe/Zurich")
V5_PROFILE_DEFAULT = "default"
V5_PROFILE_HEALTH_RECORD_6E = "health-record-6e"
ALLOWED_V5_PROFILES = frozenset({V5_PROFILE_DEFAULT, V5_PROFILE_HEALTH_RECORD_6E})
MAPPING_DECISIONS = frozenset(
    {
        "assign",
        "composite",
        "ignore",
        "defer",
        "not_assignable",
        "irrelevant",
        "conflict",
        "reopen",
    }
)
MAPPING_CONFIDENCE = frozenset({"low", "medium", "high"})
MAPPING_METHODS = frozenset({"sighi_reference", "ingredient_label", "manual_review", "local_alias"})
PERSONAL_TOLERANCE_STATUSES = frozenset(
    {"unknown", "documented_tolerated", "documented_not_tolerated", "unclear"}
)
HISTAMINE_DEFINITION_VERSION = "sighi_mapping_load_v2"


def parse_v5_profile(raw: str | None) -> str:
    """Return the one allowed regeneration profile or fail before any mutation."""
    value = V5_PROFILE_DEFAULT if raw is None or raw == "" else raw
    if value not in ALLOWED_V5_PROFILES:
        raise RuntimeError("invalid HEALTH_DASHBOARD_V5_PROFILE")
    return value


DASHBOARD_V5_PROFILE = parse_v5_profile(os.environ.get("HEALTH_DASHBOARD_V5_PROFILE"))


def local_today() -> date:
    return datetime.now(LOCAL_TIMEZONE).date()


def validate_symptom_payload(payload: dict[str, Any]) -> dict[str, Any]:
    if set(payload) != {
        "version",
        "action",
        "date",
        "scores",
        "notes",
    }:
        raise ValueError("invalid payload shape")
    if payload["version"] != 1 or payload["action"] != "symptom_checkin":
        raise ValueError("unsupported action")
    try:
        day = date.fromisoformat(payload["date"])
    except (TypeError, ValueError) as exc:
        raise ValueError("invalid date") from exc
    if day.year < 2000 or day > local_today():
        raise ValueError("date outside allowed range")
    scores = payload["scores"]
    if not isinstance(scores, dict) or set(scores) != set(FIELDS):
        raise ValueError("incomplete scores")
    if any(
        type(scores[field]) is not int or scores[field] not in range(4)
        for field in FIELDS
    ):
        raise ValueError("invalid score")
    notes = payload["notes"]
    if not isinstance(notes, str) or len(notes) > 300:
        raise ValueError("invalid notes")
    return {
        "version": 1,
        "action": "symptom_checkin",
        "date": day.isoformat(),
        "scores": {field: scores[field] for field in FIELDS},
        "notes": notes.strip(),
    }


def validate_mapping_payload(payload: dict[str, Any]) -> dict[str, Any]:
    required = {
        "version",
        "action",
        "decision",
        "queue_key",
        "alias",
        "canonical_food",
        "sighi_score",
        "confidence",
        "mapping_method",
        "source_label",
        "source_version",
        "note",
        "ingredient_review_required",
        "personal_tolerance_status",
        "personal_tolerance_note",
        "target_revision",
        "target_count",
        "target_day_count",
    }
    if set(payload) != required:
        raise ValueError("invalid mapping payload shape")
    if payload["version"] != 2 or payload["action"] != "nutrition_mapping":
        raise ValueError("unsupported action")
    decision = payload["decision"]
    queue_key = payload["queue_key"]
    alias = " ".join(str(payload["alias"] or "").split())
    canonical = " ".join(str(payload["canonical_food"] or "").split())
    score = payload["sighi_score"]
    confidence = payload["confidence"]
    method = payload["mapping_method"]
    source_label = " ".join(str(payload["source_label"] or "").split())
    source_version = " ".join(str(payload["source_version"] or "").split())
    note = " ".join(str(payload["note"] or "").split())
    ingredient_review = payload["ingredient_review_required"]
    tolerance_status = str(payload["personal_tolerance_status"] or "")
    tolerance_note = " ".join(str(payload["personal_tolerance_note"] or "").split())
    target_revision = payload["target_revision"]
    target_count = payload["target_count"]
    target_day_count = payload["target_day_count"]
    if not isinstance(queue_key, str) or not re.fullmatch(r"[0-9a-f]{32}", queue_key):
        raise ValueError("invalid queue key")
    if not isinstance(target_revision, str) or not re.fullmatch(r"[0-9a-f]{64}", target_revision):
        raise ValueError("invalid mapping target revision")
    if (
        type(target_count) is not int
        or type(target_day_count) is not int
        or not 1 <= target_count <= 10000
        or not 1 <= target_day_count <= 3660
    ):
        raise ValueError("invalid mapping target counts")
    if decision not in MAPPING_DECISIONS or confidence not in MAPPING_CONFIDENCE or method not in MAPPING_METHODS:
        raise ValueError("invalid mapping value")
    for field_name, text, maximum in (
        ("alias", alias, 160),
        ("canonical food", canonical, 120),
        ("source label", source_label, 80),
        ("source version", source_version, 40),
        ("note", note, 300),
        ("personal tolerance note", tolerance_note, 160),
    ):
        if len(text) > maximum or re.search(r"[\x00-\x1f\x7f]|https?://|/|\\\\|\.hermes", text, re.I):
            raise ValueError(f"invalid {field_name}")
    if not alias or not source_label or not source_version:
        raise ValueError("required mapping documentation missing")
    if score != "unknown" and (type(score) is not int or score not in range(4)):
        raise ValueError("invalid score")
    if type(ingredient_review) is not bool:
        raise ValueError("invalid ingredient review flag")
    if tolerance_status not in PERSONAL_TOLERANCE_STATUSES:
        raise ValueError("invalid personal tolerance status")
    if decision == "assign":
        if not canonical or score == "unknown" or ingredient_review:
            raise ValueError("assign requires canonical score without ingredient review")
    elif decision == "composite":
        if score != "unknown" or not ingredient_review:
            raise ValueError("composite must stay unclassified for ingredient review")
    elif decision == "ignore":
        if score != "unknown" or canonical or ingredient_review:
            raise ValueError("ignore must not classify")
    elif decision in {"defer", "not_assignable", "irrelevant", "conflict", "reopen"}:
        if (
            score != "unknown"
            or canonical
            or ingredient_review
            or tolerance_status != "unknown"
            or tolerance_note
        ):
            raise ValueError("review-state action must not classify")
    return {
        "version": 2,
        "action": "nutrition_mapping",
        "decision": decision,
        "queue_key": queue_key,
        "alias": alias,
        "canonical_food": canonical,
        "sighi_score": score,
        "confidence": confidence,
        "mapping_method": method,
        "source_label": source_label,
        "source_version": source_version,
        "note": note,
        "ingredient_review_required": ingredient_review,
        "personal_tolerance_status": tolerance_status,
        "personal_tolerance_note": tolerance_note,
        "target_revision": target_revision,
        "target_count": target_count,
        "target_day_count": target_day_count,
    }


def _validate_occurred_at(value: Any) -> str:
    if not isinstance(value, str):
        raise ValueError("invalid occurred_at")
    try:
        parsed = datetime.fromisoformat(value)
    except ValueError as exc:
        raise ValueError("invalid occurred_at") from exc
    if parsed.tzinfo is not None or parsed.year < 2000 or parsed.replace(tzinfo=LOCAL_TIMEZONE) > datetime.now(LOCAL_TIMEZONE):
        raise ValueError("invalid occurred_at")
    return parsed.isoformat(timespec="minutes")


def _payload_text(value: Any, maximum: int, *, required: bool = False) -> str:
    if not isinstance(value, str):
        raise ValueError("invalid text")
    cleaned = " ".join(value.split())
    if (required and not cleaned) or len(cleaned) > maximum or re.search(r"[\x00-\x1f\x7f]|https?://|/|\\\\|\.hermes", cleaned, re.I):
        raise ValueError("invalid text")
    return cleaned


def validate_patient_payload(payload: dict[str, Any]) -> dict[str, Any]:
    action = payload.get("action")
    if action == "symptom_event":
        required = {"version", "action", "occurred_at", "symptom_type", "severity", "onset_at", "duration_minutes", "label", "notes"}
        if set(payload) != required or payload.get("version") != 1 or payload.get("symptom_type") not in {"headache", "aphthae", "gi", "joints", "skin", "eyes", "fatigue", "other"} or type(payload.get("severity")) is not int or payload["severity"] not in range(4):
            raise ValueError("invalid symptom event")
        onset = _validate_occurred_at(payload["onset_at"]) if payload["onset_at"] else ""
        duration = payload["duration_minutes"]
        if duration is not None and (type(duration) is not int or duration < 0 or duration > 10080):
            raise ValueError("invalid symptom duration")
        return {"version": 1, "action": action, "occurred_at": _validate_occurred_at(payload["occurred_at"]), "symptom_type": payload["symptom_type"], "severity": payload["severity"], "onset_at": onset, "duration_minutes": duration, "label": _payload_text(payload["label"], 80, required=payload["symptom_type"] == "other"), "notes": _payload_text(payload["notes"], 300)}
    if action == "medication_event":
        required = {"version", "action", "occurred_at", "event_type", "medication_name", "dose", "unit", "route", "notes"}
        if set(payload) != required or payload.get("version") != 1 or payload.get("event_type") not in {"administered", "missed", "corrected"}:
            raise ValueError("invalid medication event")
        return {"version": 1, "action": action, "occurred_at": _validate_occurred_at(payload["occurred_at"]), "event_type": payload["event_type"], "medication_name": _payload_text(payload["medication_name"], 120, required=True), "dose": _payload_text(payload["dose"], 40), "unit": _payload_text(payload["unit"], 30), "route": _payload_text(payload["route"], 60), "notes": _payload_text(payload["notes"], 300)}
    if action == "general_event":
        required = {"version", "action", "occurred_at", "category", "label", "intensity", "notes"}
        if set(payload) != required or payload.get("version") != 1 or payload.get("category") not in {"stress", "infection", "appointment", "physical_load", "sleep_disruption", "heat", "travel", "other"}:
            raise ValueError("invalid general event")
        intensity = payload["intensity"]
        if intensity is not None and (type(intensity) is not int or intensity not in range(4)):
            raise ValueError("invalid event intensity")
        return {"version": 1, "action": action, "occurred_at": _validate_occurred_at(payload["occurred_at"]), "category": payload["category"], "label": _payload_text(payload["label"], 100, required=True), "intensity": intensity, "notes": _payload_text(payload["notes"], 300)}
    raise ValueError("unsupported patient action")


def validate_document_action(payload: dict[str, Any]) -> dict[str, Any]:
    action = payload.get("action")
    if action == "document_import":
        if set(payload) != {"version", "action", "quarantine_token", "sha256", "metadata"} or payload.get("version") != 1:
            raise ValueError("invalid document import")
        token, digest = payload["quarantine_token"], payload["sha256"]
        if not isinstance(token, str) or not TOKEN_RE.fullmatch(token) or not isinstance(digest, str) or not re.fullmatch(r"[a-f0-9]{64}", digest):
            raise ValueError("invalid document import")
        return {"version":1,"action":action,"quarantine_token":token,"sha256":digest,"metadata":validate_metadata(payload["metadata"])}
    if action != "document_review" or payload.get("version") != 1:
        raise ValueError("unsupported document action")
    required = {"version","action","operation","document_id","target_id","value","unit","decision","metadata"}
    optional = {"expected_candidate_revision", "preview_revision"}
    if not required.issubset(payload) or not set(payload).issubset(required | optional) or not isinstance(payload["document_id"], str) or not ID_RE.fullmatch(payload["document_id"]):
        raise ValueError("invalid document review")
    operation = payload["operation"]
    if 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 review")
    target = payload["target_id"]
    if not isinstance(target, str) or len(target) > 64:
        raise ValueError("invalid target")
    if operation in {"candidate_decision","candidate_remap","candidate_transfer"} and not CANDIDATE_ID_RE.fullmatch(target):
        raise ValueError("invalid candidate")
    if operation in {"text_correction","page_review"} and not re.fullmatch(r"page_[1-9][0-9]{0,2}", target):
        raise ValueError("invalid page target")
    def review_text(raw: Any, maximum: int, *, required: bool = False) -> str:
        if not isinstance(raw, str): raise ValueError("invalid review text")
        clean = " ".join(raw.split())
        if (required and not clean) or len(clean) > maximum or re.search(r"[\x00-\x1f\x7f]|https?://|\\\\|\.hermes", clean, re.I):
            raise ValueError("invalid review text")
        return clean
    value = review_text(payload["value"], 4000, required=operation in {"text_correction","candidate_decision","candidate_remap"})
    unit = review_text(payload["unit"], 40)
    decision = review_text(payload["decision"], 30, required=operation == "candidate_decision")
    expected_revision = payload.get("expected_candidate_revision")
    preview_revision = payload.get("preview_revision")
    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")
    metadata = validate_metadata(payload["metadata"]) if operation == "metadata_review" else validate_metadata({"document_date":"","document_type":"other","institution":"","personal_title":"","investigation_day":"","note":""})
    result = {"version":1,"action":action,"operation":operation,"document_id":payload["document_id"],"target_id":target,"value":value,"unit":unit,"decision":decision,"metadata":metadata}
    if expected_revision is not None: result["expected_candidate_revision"] = expected_revision
    if preview_revision is not None: result["preview_revision"] = preview_revision
    return result


def validate_payload(payload: Any) -> dict[str, Any]:
    if not isinstance(payload, dict):
        raise ValueError("invalid payload shape")
    action = payload.get("action")
    if action in {"document_import", "document_review"}:
        return validate_document_action(payload)
    if action == "symptom_checkin":
        return validate_symptom_payload(payload)
    if action == "nutrition_mapping":
        return validate_mapping_payload(payload)
    if action == "capture_entry":
        return validate_capture_payload(payload)
    if action in {"symptom_event", "medication_event", "general_event"}:
        return validate_patient_payload(payload)
    if action in {"supplement_plan", "supplement_intake"}:
        return validate_supplement_payload(payload)
    if action in {
        "observation_upsert",
        "observation_phase_upsert",
        "observation_checkin",
        "observation_status",
        "observation_result_snapshot",
    }:
        return validate_observation_action(payload)
    raise ValueError("unsupported action")


def reject_duplicate_object_pairs(pairs: list[tuple[str, object]]) -> dict[str, object]:
    result: dict[str, object] = {}
    for key, value in pairs:
        if key in result:
            raise ValueError("duplicate JSON key")
        result[key] = value
    return result


def ensure_private_inbox() -> None:
    ACTION_INBOX.mkdir(mode=0o700, parents=True, exist_ok=True)
    metadata = ACTION_INBOX.lstat()
    if not stat.S_ISDIR(metadata.st_mode) or metadata.st_uid != os.getuid():
        raise OSError("action inbox is not a private owned directory")
    os.chmod(ACTION_INBOX, 0o700)


def load_action(path: Path) -> dict[str, Any]:
    if not NAME_RE.fullmatch(path.name):
        raise ValueError("invalid filename")
    descriptor = os.open(path, os.O_RDONLY | os.O_NOFOLLOW)
    try:
        metadata = os.fstat(descriptor)
        if (
            not stat.S_ISREG(metadata.st_mode)
            or metadata.st_uid != os.getuid()
            or stat.S_IMODE(metadata.st_mode) != 0o600
        ):
            raise ValueError("invalid action file permissions")
        if metadata.st_size > 8192:
            raise ValueError("invalid action file")
        raw = os.read(descriptor, 8193)
    finally:
        os.close(descriptor)
    if len(raw) > 8192:
        raise ValueError("action too large")
    return validate_payload(
        json.loads(raw.decode("utf-8"), object_pairs_hook=reject_duplicate_object_pairs)
    )


def apply_symptom_action(day: str, scores: dict[str, int], notes: str) -> None:
    if DASHBOARD_DB is None:
        raise RuntimeError("HEALTH_DASHBOARD_DB is required")
    command = [
        sys.executable,
        str(QUICK_ADD),
        "--db",
        str(DASHBOARD_DB),
        "--date",
        day,
        "--notes",
        notes,
    ]
    for field in FIELDS:
        command.extend((f"--{field}", str(scores[field])))
    subprocess.run(command, check=True, timeout=120, capture_output=True, text=True)
    if DASHBOARD_DB == DEFAULT_DASHBOARD_DB:
        subprocess.run(
            [sys.executable, str(DASHBOARD_V4), "--output", str(DASHBOARD_FILE)],
            check=True,
            timeout=120,
            capture_output=True,
            text=True,
        )


def nutrition_queue_key(normalized_name: str) -> str:
    material = b"health-nutrition-queue-v1:" + normalized_name.encode("utf-8")
    return hashlib.sha256(material).hexdigest()[:32]


def _score_label(max_score: int, unknown: int, load: float) -> str:
    if unknown:
        return "unknown"
    return "classified"


def _public_numeric(value: Any) -> float | None:
    try:
        numeric = float(value)
    except (TypeError, ValueError):
        return None
    return numeric if numeric == numeric and numeric not in {float("inf"), float("-inf")} else None


NUTRIENT_ALLOWLIST = {
    "energy.energy": "kcal",
    "nutrient.protein": "g",
    "nutrient.carb": "g",
    "nutrient.fat": "g",
    "nutrient.fiber": "g",
    "nutrient.sugar": "g",
    "nutrient.saturated": "g",
    "nutrient.salt": "g",
    "nutrient.sodium": "mg",
}


def _nutrient_totals(connection: sqlite3.Connection, item_ids: list[int]) -> str:
    if not item_ids:
        return "{}"
    placeholders = ",".join("?" for _ in item_ids)
    totals: dict[str, float] = {}
    for row in connection.execute(
        f"""SELECT nutrient_key,value,unit FROM nutrition_item_nutrients
            WHERE item_id IN ({placeholders})""",
        item_ids,
    ):
        key = str(row["nutrient_key"] or "")
        unit = str(row["unit"] or "")
        value = _public_numeric(row["value"])
        if key not in NUTRIENT_ALLOWLIST or unit != NUTRIENT_ALLOWLIST[key] or value is None:
            continue
        totals[key] = totals.get(key, 0.0) + value
    return json.dumps({key: round(value, 3) for key, value in sorted(totals.items())}, separators=(",", ":"))


def assert_nutrition_mapping_schema(connection: sqlite3.Connection) -> None:
    """Fail closed without creating or altering schema in the action path."""
    required = {
        "nutrition_items": {"id", "item_hash", "datum", "name"},
        "nutrition_histamine_scores": {"item_id", "canonical_food", "sighi_score"},
        "nutrition_review_queue": {"normalized_name", "example_name", "status"},
        "histamine_food_rules": {"canonical_food", "sighi_score", "confidence", "source", "updated_at"},
        "histamine_food_aliases": {"alias", "canonical_food"},
        "nutrition_mapping_action_log": {"action_id", "action_hash", "queue_key", "decision"},
        "nutrition_mapping_provenance": {"action_id", "normalized_name", "affected_day"},
        "nutrition_composite_product_review": {"normalized_name", "status"},
        "personal_food_tolerance": {"canonical_food", "personal_status"},
        "nutrition_meal_summary": {"datum", "meal", "histamine_score", "histamine_label"},
    }
    for table, expected in required.items():
        exists = connection.execute(
            "SELECT 1 FROM sqlite_master WHERE type='table' AND name=?", (table,)
        ).fetchone()
        if exists is None:
            raise RuntimeError("nutrition mapping schema missing")
        actual = {
            str(row["name"] if isinstance(row, sqlite3.Row) else row[1])
            for row in connection.execute(f'PRAGMA table_info("{table}")')
        }
        if not expected <= actual:
            raise RuntimeError("nutrition mapping schema missing")


def _classified_score(row: sqlite3.Row) -> int | None:
    score = row["sighi_score"]
    canonical = str(row["canonical_food"] or "").strip()
    traffic = str(row["traffic_light"] or "").strip()
    confidence = str(row["confidence"] or "").strip()
    if (
        type(score) is not int
        or not 0 <= score <= 3
        or not canonical
        or traffic not in {"classified", "green", "yellow", "orange", "red"}
        or confidence not in MAPPING_CONFIDENCE
    ):
        return None
    return score


def recompute_nutrition_day(connection: sqlite3.Connection, day: str) -> None:
    rows = list(
        connection.execute(
            """SELECT i.id,i.meal,i.kcal,i.protein_g,i.carb_g,i.fat_g,
                      h.canonical_food,h.sighi_score,h.traffic_light,h.confidence
               FROM nutrition_items i LEFT JOIN nutrition_histamine_scores h ON h.item_id=i.id
               WHERE i.datum=?""",
            (day,),
        )
    )
    connection.execute("DELETE FROM nutrition_meal_summary WHERE datum=?", (day,))
    if not rows:
        connection.execute("DELETE FROM nutrition_daily_summary_v2 WHERE datum=?", (day,))
        return
    totals = {"kcal": 0.0, "protein": 0.0, "carb": 0.0, "fat": 0.0, "items": 0, "unknown": 0, "load": 0.0, "scores": [], "missing": set(), "ids": []}
    by_meal: dict[str, dict[str, Any]] = {}
    for row in rows:
        meal = row["meal"] if row["meal"] in {"breakfast", "lunch", "dinner", "snack"} else "unassigned"
        bucket = by_meal.setdefault(meal, {"kcal": 0.0, "protein": 0.0, "carb": 0.0, "fat": 0.0, "items": 0, "unknown": 0, "load": 0.0, "scores": [], "missing": set(), "ids": []})
        for target in (bucket, totals):
            target["items"] += 1
            target["ids"].append(int(row["id"]))
            for source, dest in (("kcal", "kcal"), ("protein_g", "protein"), ("carb_g", "carb"), ("fat_g", "fat")):
                if row[source] is None:
                    target["missing"].add(dest)
                else:
                    target[dest] += float(row[source])
            score = _classified_score(row)
            if score is None:
                target["unknown"] += 1
            else:
                target["scores"].append(score)
                target["load"] += score
    for meal, bucket in by_meal.items():
        max_score = max(bucket["scores"] or [0])
        label = _score_label(max_score, int(bucket["unknown"]), float(bucket["load"]))
        values = {key: None if key in bucket["missing"] else round(float(bucket[key]), 1) for key in ("kcal", "protein", "carb", "fat")}
        connection.execute(
            """INSERT OR REPLACE INTO nutrition_meal_summary
               (datum,meal,kcal,protein_g,carb_g,fat_g,item_count,nutrient_json,histamine_score,histamine_label,updated_at)
               VALUES(?,?,?,?,?,?,?,?,?,?,CURRENT_TIMESTAMP)""",
            (day, meal, values["kcal"], values["protein"], values["carb"], values["fat"], int(bucket["items"]), _nutrient_totals(connection, bucket["ids"]), round(bucket["load"], 1) if not bucket["unknown"] else None, label),
        )
    max_score = max(totals["scores"] or [0])
    label = _score_label(max_score, int(totals["unknown"]), float(totals["load"]))
    connection.execute(
        """INSERT OR REPLACE INTO nutrition_daily_summary_v2
           (datum,kcal,protein_g,carb_g,fat_g,item_count,nutrient_json,histamine_score,histamine_max,histamine_unknown_count,histamine_label,updated_at)
           VALUES(?,?,?,?,?,?,?,?,?,?,?,CURRENT_TIMESTAMP)""",
        (
            day,
            None if "kcal" in totals["missing"] else round(float(totals["kcal"]), 1),
            None if "protein" in totals["missing"] else round(float(totals["protein"]), 1),
            None if "carb" in totals["missing"] else round(float(totals["carb"]), 1),
            None if "fat" in totals["missing"] else round(float(totals["fat"]), 1),
            int(totals["items"]),
            _nutrient_totals(connection, totals["ids"]),
            round(totals["load"], 1) if not totals["unknown"] else None,
            max_score,
            int(totals["unknown"]),
            label,
        ),
    )


def apply_mapping_action(payload: dict[str, Any], action_id: str) -> None:
    """Apply one preview-bound group review atomically and idempotently."""
    if DASHBOARD_DB is None:
        raise RuntimeError("HEALTH_DASHBOARD_DB is required")
    action_hash = hashlib.sha256(
        json.dumps(payload, sort_keys=True, separators=(",", ":")).encode("utf-8")
    ).hexdigest()
    connection = sqlite3.connect(DASHBOARD_DB)
    connection.row_factory = sqlite3.Row
    try:
        connection.execute("PRAGMA foreign_keys=ON")
        assert_nutrition_mapping_schema(connection)
        connection.execute("BEGIN IMMEDIATE")
        with connection:
            existing = connection.execute(
                "SELECT action_hash FROM nutrition_mapping_action_log WHERE action_id=?",
                (action_id,),
            ).fetchone()
            if existing is not None:
                if existing["action_hash"] != action_hash:
                    raise RuntimeError("conflicting idempotency key")
                return

            group = resolve_mapping_group(connection, payload["queue_key"])
            if group is None:
                raise RuntimeError("mapping group not found")
            if (
                group["target_revision"] != payload["target_revision"]
                or int(group["target_count"]) != payload["target_count"]
                or int(group["target_day_count"]) != payload["target_day_count"]
            ):
                raise RuntimeError("mapping preview target changed")
            example_name = str(group["example_name"] or "")
            if _exact_food_identity(payload["alias"]) != _exact_food_identity(example_name):
                raise RuntimeError("mapping alias does not match group identity")
            normalized = _normalize_food_identity(example_name)
            item_cursor = connection.execute(
                    """SELECT i.id,i.item_hash,i.datum,i.name,
                              h.item_id AS mapped_item_id,h.canonical_food,h.sighi_score,
                              h.traffic_light,h.tags,h.confidence,h.reason
                         FROM nutrition_items i
                         LEFT JOIN nutrition_histamine_scores h ON h.item_id=i.id"""
                )
            all_items = item_cursor.fetchmany(250_001)
            if len(all_items) > 250_000:
                raise RuntimeError("nutrition mapping row limit exceeded")
            item_rows = [
                row for row in all_items
                if _exact_food_identity(row["name"]) == _exact_food_identity(example_name)
            ]
            if len(item_rows) != payload["target_count"]:
                raise RuntimeError("mapping preview target changed")
            affected = sorted({str(row["datum"]) for row in item_rows if row["datum"]})
            if len(affected) != payload["target_day_count"] or len(affected) > 3660:
                raise RuntimeError("mapping preview target changed")

            connection.execute(
                """INSERT OR IGNORE INTO nutrition_review_queue
                   (normalized_name,example_name,occurrence_count,first_seen,last_seen,
                    suggested_canonical_food,suggested_score,reason,status,updated_at)
                   VALUES(?,?,?,?,?,NULL,NULL,'derived exact review group','open',CURRENT_TIMESTAMP)""",
                (normalized, example_name, len(item_rows), affected[0], affected[-1]),
            )
            canonical = payload["canonical_food"]
            score = None if payload["sighi_score"] == "unknown" else payload["sighi_score"]
            effective_source_label = payload["source_label"]
            effective_source_version = payload["source_version"]
            effective_method = payload["mapping_method"]
            effective_confidence = payload["confidence"]
            if payload["decision"] == "assign":
                require_alias = payload["mapping_method"] == "local_alias"
                rule = resolve_catalog_assignment(
                    connection,
                    canonical,
                    alias=example_name,
                    require_alias=require_alias,
                )
                if rule is None or rule["sighi_score"] != score:
                    raise RuntimeError("local reference is absent, ambiguous, or version-mismatched")
                effective_method = "local_alias" if require_alias else "sighi_reference"
                effective_confidence = rule["confidence"]
                effective_source_label = rule["provenance_source"]
                effective_source_version = rule["provenance_version"]

            decision = payload["decision"]
            effective_definition_version = HISTAMINE_DEFINITION_VERSION
            if decision == "assign":
                if score is None:
                    raise RuntimeError("assigned mapping score is missing")
                effective_definition_version += "|" + assignment_evidence_token(
                    item_rows, canonical, score
                )
            if decision == "composite":
                connection.execute(
                    """INSERT OR REPLACE INTO nutrition_composite_product_review
                       (normalized_name,canonical_food,status,updated_at)
                       VALUES(?,?, 'needs_ingredient_review', CURRENT_TIMESTAMP)""",
                    (normalized, None),
                )
            elif decision == "assign":
                for row in item_rows:
                    connection.execute(
                        """INSERT INTO nutrition_histamine_scores
                           (item_id,canonical_food,sighi_score,traffic_light,tags,confidence,reason,scored_at)
                           VALUES(?,?,?,?,?,?,?,CURRENT_TIMESTAMP)
                           ON CONFLICT(item_id) DO UPDATE SET
                             canonical_food=excluded.canonical_food,
                             sighi_score=excluded.sighi_score,
                             traffic_light=excluded.traffic_light,
                             tags=excluded.tags,
                             confidence=excluded.confidence,
                             reason=excluded.reason,
                             scored_at=CURRENT_TIMESTAMP""",
                        (
                            row["id"], canonical, score, "classified", effective_method,
                            effective_confidence, payload["note"] or "user-confirmed mapping review",
                        ),
                    )
                if payload["personal_tolerance_status"] != "unknown":
                    connection.execute(
                        """INSERT INTO personal_food_tolerance
                           (canonical_food,personal_status,evidence_level,notes,updated_at)
                           VALUES(?,?, 'user_documented', ?, CURRENT_TIMESTAMP)
                           ON CONFLICT(canonical_food) DO UPDATE SET
                             personal_status=excluded.personal_status,
                             evidence_level=excluded.evidence_level,
                             notes=excluded.notes,
                             updated_at=CURRENT_TIMESTAMP""",
                        (
                            canonical,
                            payload["personal_tolerance_status"],
                            payload["personal_tolerance_note"] or None,
                        ),
                    )
                unresolved_same_legacy = any(
                    _normalize_food_identity(str(row["name"] or "")) == normalized
                    and _exact_food_identity(row["name"]) != _exact_food_identity(example_name)
                    and _classified_score(row) is None
                    for row in all_items
                )
                connection.execute(
                    "UPDATE nutrition_review_queue SET status=?,updated_at=CURRENT_TIMESTAMP WHERE normalized_name=?",
                    ("open" if unresolved_same_legacy else "mapped", normalized),
                )
            elif decision in {"defer", "conflict", "reopen"}:
                connection.execute(
                    "UPDATE nutrition_review_queue SET status='open',updated_at=CURRENT_TIMESTAMP WHERE normalized_name=?",
                    (normalized,),
                )
            elif decision in {"not_assignable", "irrelevant", "ignore"}:
                # The action log is the fine-grained review truth.  The legacy queue
                # stays open if another exact identity shares its aggressive key.
                sibling = any(
                    _normalize_food_identity(str(row["name"] or "")) == normalized
                    and _exact_food_identity(row["name"]) != _exact_food_identity(example_name)
                    for row in all_items
                )
                connection.execute(
                    "UPDATE nutrition_review_queue SET status=?,updated_at=CURRENT_TIMESTAMP WHERE normalized_name=?",
                    ("open" if sibling else "ignored", normalized),
                )

            for day in affected:
                if decision == "assign":
                    recompute_nutrition_day(connection, day)
                connection.execute(
                    """INSERT OR IGNORE INTO nutrition_mapping_provenance
                       (action_id,normalized_name,decision,alias,canonical_food,sighi_score,confidence,
                        mapping_method,source_label,source_version,note,ingredient_review_required,
                        personal_tolerance_status,personal_tolerance_note,definition_version,affected_day)
                       VALUES(?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)""",
                    (
                        action_id, normalized, decision, payload["alias"], canonical or None,
                        score, effective_confidence, effective_method, effective_source_label,
                        effective_source_version, payload["note"], int(payload["ingredient_review_required"]),
                        payload["personal_tolerance_status"], payload["personal_tolerance_note"] or None,
                        effective_definition_version, day,
                    ),
                )
            connection.execute(
                """INSERT INTO nutrition_mapping_action_log
                   (action_id,action_hash,queue_key,decision,alias,canonical_food,sighi_score,confidence,
                    mapping_method,source_label,source_version,note,ingredient_review_required,
                    personal_tolerance_status,personal_tolerance_note,definition_version)
                   VALUES(?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)""",
                (
                    action_id, action_hash, payload["queue_key"], decision, payload["alias"],
                    canonical or None, score, effective_confidence, effective_method,
                    effective_source_label, effective_source_version, payload["note"],
                    int(payload["ingredient_review_required"]), payload["personal_tolerance_status"],
                    payload["personal_tolerance_note"] or None, effective_definition_version,
                ),
            )
    finally:
        connection.close()


def _assert_patient_action_schema(connection: sqlite3.Connection) -> None:
    for table, columns in REQUIRED_PATIENT_ACTION_COLUMNS.items():
        exists = connection.execute(
            "SELECT 1 FROM sqlite_master WHERE type='table' AND name=?", (table,)
        ).fetchone()
        if exists is None:
            raise RuntimeError("patient-action schema missing")
        actual = {
            str(row[1]) for row in connection.execute(f'PRAGMA table_info("{table}")')
        }
        if any(name not in actual for name, _kind in columns):
            raise RuntimeError("patient-action schema missing")


def apply_supplement_action(payload: dict[str, Any]) -> None:
    if DASHBOARD_DB is None:
        raise RuntimeError("HEALTH_DASHBOARD_DB is required")
    connection = sqlite3.connect(DASHBOARD_DB)
    try:
        connection.execute("PRAGMA foreign_keys=ON")
        assert_supplement_schema(connection)
        connection.execute("BEGIN IMMEDIATE")
        if payload["action"] == "supplement_plan":
            connection.execute(
                """INSERT INTO supplement_plans
                   (product,brand_variant,nutrient_key,amount,unit,schedule_type,weekdays,
                    interval_days,start_date,end_date,composition_source,
                    assignment_reliability,notes)
                   VALUES(?,?,?,?,?,?,?,?,?,?,?,?,?)""",
                (
                    payload["product"], payload["brand_variant"] or None,
                    payload["nutrient_key"], payload["amount"], payload["unit"],
                    payload["schedule_type"], payload["weekdays"] or None,
                    payload["interval_days"], payload["start_date"],
                    payload["end_date"] or None, payload["composition_source"],
                    payload["assignment_reliability"], payload["notes"] or None,
                ),
            )
        else:
            connection.execute(
                """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") 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")
            if data["contract"] == "health.medication_action.v2" and data["mode"] == "correction":
                if correction_target is None:
                    raise RuntimeError("structured correction target missing")
                linked_capture = connection.execute(
                    "SELECT entry_id FROM capture_action_log WHERE action_hash=?",
                    (str(correction_target["business_revision"] or ""),),
                ).fetchone()
                if linked_capture is None or str(linked_capture[0]) != payload["corrects_entry_id"]:
                    raise RuntimeError("structured correction capture target changed")
            medication_id = int(prescription["id"])
            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,
                "planned_quantity_value": str(planned["planned_quantity_value"] or "") if planned is not None else "",
                "planned_dosage_form": str(planned["planned_dosage_form"] or "") if planned is not None else "",
                "planned_strength": str(planned["planned_strength"] or "") if planned is not None else "",
                "corrects_event_id": int(correction_target["id"]) if correction_target is not None else None,
                "prescription_business_revision": str(prescription["business_revision"] or ""),
            }
        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"
            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:
                    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,planned_quantity_value,planned_dosage_form,
                            planned_strength,route_original,route_normalized,injection_region,
                            injection_side,injection_detail,business_revision,actual_quantity_value,
                            actual_dosage_form,actual_strength,corrects_event_id,corrected_target_status,
                            correction_reason)
                            VALUES(?,?,?,?,?,NULL,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)""",
                            (
                                day, data["name"], legacy_dose, legacy_route or None,
                                data["status"], data["note"] or None,
                                "dashboard_medication_action_version_two", payload["occurred_at"],
                                medication_context["medication_id"], medication_context["planned_event_id"],
                                medication_context["planned_quantity_value"] or None,
                                medication_context["planned_dosage_form"] or None,
                                medication_context["planned_strength"] 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, action_hash,
                                data["quantity_value"], data["dosage_form"], data["strength"],
                                medication_context["corrects_event_id"],
                                "administered" if data["mode"] == "correction" else None,
                                data.get("correction_reason") or None,
                            ),
                        )
                        if data.get("bind_verified_preset"):
                            if (
                                data["mode"] != "correction"
                                or (
                                    data["quantity_value"], data["dosage_form"],
                                    data["strength"], data["route_normalized"],
                                ) != ("1", "Spritze", "40 mg/0,4 ml", "subcutaneous")
                            ):
                                raise RuntimeError("verified preset confirmation changed")
                            bind_verified_hyrimoz_preset(
                                connection,
                                medication_context["medication_id"],
                                str(medication_context["prescription_business_revision"]),
                            )
                    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,
                            value,
                            unit,
                            source,
                            data["note"] or None,
                        ),
                    )
            elif payload["capture_type"] == "measurement" and "health_events" in available_tables:
                event_columns = {str(row[1]) for row in connection.execute("PRAGMA table_info(health_events)")}
                for key, value in data["values"].items():
                    if "occurred_at" in event_columns:
                        connection.execute(
                            "INSERT INTO health_events(date,category,parameter,value,unit,source,notes,occurred_at,intensity) VALUES(?,'manual_measurement',?,?,?,?,?,?,NULL)",
                            (day, key, value, data["unit"], "dashboard_v5_mobile_capture", data["note"] or None, payload["occurred_at"]),
                        )
                    else:
                        connection.execute(
                            "INSERT INTO health_events(date,category,parameter,value,unit,source,notes) VALUES(?,'manual_measurement',?,?,?,?,?)",
                            (day, key, value, data["unit"], "dashboard_v5_mobile_capture", data["note"] or None),
                        )
        connection.execute(
            "INSERT INTO capture_action_log(idempotency_key,request_version,action_hash,entry_id,processed_at) VALUES(?,?,?,?,?)",
            (
                payload["idempotency_key"],
                payload["request_version"],
                action_hash,
                entry_id,
                now,
            ),
        )
        connection.commit()
        cleanup_expired(CAPTURE_QUARANTINE)
        return entry_id
    except Exception:
        connection.rollback()
        for attachment in promoted:
            if attachment.get("reused"):
                continue
            for name in (
                attachment.get("media_name"),
                attachment.get("thumbnail_name"),
                attachment.get("proxy_name"),
            ):
                if name:
                    (CAPTURE_MEDIA / name).unlink(missing_ok=True)
        raise
    finally:
        connection.close()

def apply_patient_action(payload: dict[str, Any]) -> None:
    if DASHBOARD_DB is None:
        raise RuntimeError("HEALTH_DASHBOARD_DB is required")
    connection = sqlite3.connect(DASHBOARD_DB)
    try:
        _assert_patient_action_schema(connection)
        connection.execute("PRAGMA foreign_keys=ON")
        connection.execute("BEGIN IMMEDIATE")
        occurred_at = payload["occurred_at"]
        day = str(occurred_at).split("T", 1)[0]
        if payload["action"] == "symptom_event":
            symptom = payload["label"] if payload["symptom_type"] == "other" else payload["symptom_type"]
            connection.execute(
                """INSERT INTO symptom_log
                   (datum,symptom,schwergrad,kontext,notizen,occurred_at,onset_at,duration_minutes)
                   VALUES(?,?,?,'additional_symptom',?,?,?,?)""",
                (day, symptom, payload["severity"], payload["notes"], occurred_at, payload["onset_at"] or None, payload["duration_minutes"]),
            )
        elif payload["action"] == "medication_event":
            known = {str(row[0]) for row in connection.execute("SELECT DISTINCT medication_name FROM medication_administrations WHERE trim(COALESCE(medication_name,''))<>''")}
            if payload["medication_name"] not in known:
                raise RuntimeError("medication is not in the known exact catalog")
            dose = " ".join(part for part in (payload["dose"], payload["unit"]) if part)
            connection.execute(
                """INSERT INTO medication_administrations
                   (datum,medication_name,dose,route,event_type,notes,source,occurred_at)
                   VALUES(?,?,?,?,?,?, 'dashboard_v5_patient_action', ?)""",
                (day, payload["medication_name"], dose or None, payload["route"] or None, payload["event_type"], payload["notes"] or None, occurred_at),
            )
        else:
            connection.execute(
                """INSERT INTO health_events
                   (date,category,parameter,value,unit,source,notes,occurred_at,intensity)
                   VALUES(?,?,?,?,NULL,'dashboard_v5_patient_action',?,?,?)""",
                (day, payload["category"], payload["label"], payload["intensity"], payload["notes"] or None, occurred_at, payload["intensity"]),
            )
        connection.commit()
    except Exception:
        connection.rollback()
        raise
    finally:
        connection.close()


def _opaque(prefix: str) -> str:
    return prefix + secrets.token_hex(12)


def _canonical_json(value: Any) -> str:
    return json.dumps(value, ensure_ascii=True, sort_keys=True, separators=(",", ":"))


def _insert_observation_snapshot(
    connection: sqlite3.Connection, observation_id: str, action_id: str
) -> None:
    from dashboard_v5.observation_engine import evaluate_observation, evaluate_plan, observation_detail, snapshot_payload
    from dashboard_v5.read_api import APIError, _series

    def loader(metric_id: str, start: date, end: date) -> dict[str, Any]:
        try:
            return _series(connection, {"metric": metric_id, "from": start.isoformat(), "to": end.isoformat(), "resolution": "day"})
        except APIError as error:
            if error.code in {"metric_not_found", "metric_not_released"}:
                return {"points": [], "unit": "", "source": {"type": "not_documented"}}
            raise

    result_id = f"result_{action_id[:24]}"
    if connection.execute(
        "SELECT 1 FROM personal_observation_results WHERE id=?", (result_id,)
    ).fetchone() is not None:
        return
    detail = observation_detail(connection, observation_id)
    if detail["influences"][0] in {"event.sauna", "event.training", "nutrition.profile", "nutrition.histamine"} and all(
        not item.startswith("event.") for item in detail["outcomes"]
    ):
        analysis = evaluate_plan(connection, observation_id, loader)
    else:
        analysis = evaluate_observation(connection, observation_id, loader)
    snapshot = snapshot_payload(analysis)
    version = int(connection.execute("SELECT COALESCE(MAX(result_version),0)+1 FROM personal_observation_results WHERE observation_id=?", (observation_id,)).fetchone()[0])
    connection.execute(
        """INSERT INTO personal_observation_results
           (id,observation_id,result_version,configuration_json,phases_json,summary_json,
            method_version,configuration_hash,created_at) VALUES(?,?,?,?,?,?,?,?,?)""",
        (result_id, observation_id, version, _canonical_json(snapshot["configuration"]),
         _canonical_json(snapshot["phases"]), _canonical_json(snapshot["summary"]),
         METHOD_VERSION, snapshot["configuration_hash"], datetime.now(LOCAL_TIMEZONE).isoformat(timespec="seconds")),
    )


def apply_observation_action(
    payload: dict[str, Any], action_id: str | None = None
) -> None:
    if DASHBOARD_DB is None:
        raise RuntimeError("HEALTH_DASHBOARD_DB is required")
    normalized_action_id = action_id or hashlib.sha256(
        _canonical_json(payload).encode()
    ).hexdigest()[:32]
    if not re.fullmatch(r"[0-9a-f]{32}", normalized_action_id):
        raise ValueError("invalid observation action id")
    connection = sqlite3.connect(DASHBOARD_DB)
    connection.row_factory = sqlite3.Row
    try:
        connection.execute("PRAGMA foreign_keys=ON")
        assert_observation_schema(connection)
        connection.execute("BEGIN IMMEDIATE")
        now = datetime.now(LOCAL_TIMEZONE).isoformat(timespec="seconds")
        action = payload["action"]
        if action == "observation_upsert":
            observation_id = payload["observation_id"] or f"obs_{normalized_action_id[:24]}"
            existing = connection.execute("SELECT status FROM personal_observations WHERE id=?", (observation_id,)).fetchone()
            if payload["status"] == "active":
                active_count = int(connection.execute(
                    "SELECT count(*) FROM personal_observations WHERE status='active' AND id<>?",
                    (observation_id,),
                ).fetchone()[0])
                if active_count >= 3:
                    raise RuntimeError("maximum active observations reached")
            if existing is None:
                if payload["observation_id"]:
                    raise RuntimeError("observation not found")
                if payload["status"] not in {"draft", "active"}:
                    raise RuntimeError("initial observation status not allowed")
                connection.execute(
                    """INSERT INTO personal_observations
                       (id,contract_version,title,question,influences_json,outcomes_json,lag_min,lag_max,
                        start_date,end_date,status,note,include_doctor,method_version,created_at,updated_at)
                       VALUES(?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)""",
                    (observation_id,1,payload["title"],payload["question"],_canonical_json(payload["influences"]),
                     _canonical_json(payload["outcomes"]),payload["lag_min"],payload["lag_max"],payload["start_date"],
                     payload["end_date"],payload["status"],payload["note"],int(payload["include_doctor"]),METHOD_VERSION,now,now),
                )
            elif payload["observation_id"] is None:
                # The queue filename is the stable command identity. Replaying
                # the same create command therefore resolves to this same row.
                pass
            else:
                old_status = str(existing["status"])
                if payload["status"] not in TRANSITIONS[old_status]:
                    raise RuntimeError("status transition not allowed")
                outside = connection.execute("SELECT 1 FROM personal_observation_phases WHERE observation_id=? AND (start_date<? OR end_date>?) LIMIT 1", (observation_id,payload["start_date"],payload["end_date"])).fetchone()
                if outside is not None:
                    raise RuntimeError("observation range excludes phase")
                connection.execute(
                    """UPDATE personal_observations SET title=?,question=?,influences_json=?,outcomes_json=?,lag_min=?,lag_max=?,
                       start_date=?,end_date=?,status=?,note=?,include_doctor=?,method_version=?,updated_at=? WHERE id=?""",
                    (payload["title"],payload["question"],_canonical_json(payload["influences"]),_canonical_json(payload["outcomes"]),
                     payload["lag_min"],payload["lag_max"],payload["start_date"],payload["end_date"],payload["status"],payload["note"],
                     int(payload["include_doctor"]),METHOD_VERSION,now,observation_id),
                )
                if old_status != "completed" and payload["status"] == "completed":
                    _insert_observation_snapshot(connection, observation_id, normalized_action_id)
        elif action == "observation_phase_upsert":
            owner = connection.execute("SELECT status,start_date,end_date FROM personal_observations WHERE id=?", (payload["observation_id"],)).fetchone()
            if owner is None or owner["status"] in {"completed","archived"}:
                raise RuntimeError("observation cannot be edited")
            if payload["start_date"] < owner["start_date"] or payload["end_date"] > owner["end_date"]:
                raise RuntimeError("phase outside observation")
            phase_id = payload["phase_id"] or _opaque("phase_")
            existing = connection.execute("SELECT observation_id FROM personal_observation_phases WHERE id=?", (phase_id,)).fetchone()
            values=(payload["phase_type"],payload["name"],payload["start_date"],payload["end_date"],payload["behavior_goal"],_canonical_json(payload["metrics"]),_canonical_json(payload["events"]),payload["lag_min"],payload["lag_max"],payload["adherence"],payload["notes"])
            if existing is None:
                if payload["phase_id"]:
                    raise RuntimeError("phase not found")
                connection.execute("""INSERT INTO personal_observation_phases
                  (id,observation_id,phase_type,name,start_date,end_date,behavior_goal,metrics_json,events_json,lag_min,lag_max,adherence,notes,created_at,updated_at)
                  VALUES(?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)""", (phase_id,payload["observation_id"],*values,now,now))
            else:
                if existing["observation_id"] != payload["observation_id"]:
                    raise RuntimeError("phase owner mismatch")
                connection.execute("""UPDATE personal_observation_phases SET phase_type=?,name=?,start_date=?,end_date=?,behavior_goal=?,metrics_json=?,events_json=?,lag_min=?,lag_max=?,adherence=?,notes=?,updated_at=? WHERE id=?""", (*values,now,phase_id))
        elif action == "observation_checkin":
            owner = connection.execute("SELECT status,start_date,end_date FROM personal_observations WHERE id=?", (payload["observation_id"],)).fetchone()
            if owner is None or owner["status"] != "active" or not owner["start_date"] <= payload["day"] <= owner["end_date"]:
                raise RuntimeError("no active observation for day")
            phase_id = payload["phase_id"] or None
            if phase_id:
                phase = connection.execute("SELECT observation_id,start_date,end_date FROM personal_observation_phases WHERE id=?", (phase_id,)).fetchone()
                if phase is None or phase["observation_id"] != payload["observation_id"] or not phase["start_date"] <= payload["day"] <= phase["end_date"]:
                    raise RuntimeError("invalid checkin phase")
            else:
                matches=connection.execute("SELECT id FROM personal_observation_phases WHERE observation_id=? AND start_date<=? AND end_date>=? ORDER BY start_date,id LIMIT 2", (payload["observation_id"],payload["day"],payload["day"])).fetchall()
                phase_id=str(matches[0]["id"]) if len(matches)==1 else None
            connection.execute("""INSERT INTO personal_observation_checkins
              (id,observation_id,phase_id,day,adherence,stress,sleep_disruption,infection,unusual_activity,travel,
               medication_change_fact,supplement_change_fact,note,created_at,updated_at)
              VALUES(?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)
              ON CONFLICT(observation_id,day) DO UPDATE SET phase_id=excluded.phase_id,adherence=excluded.adherence,
               stress=excluded.stress,sleep_disruption=excluded.sleep_disruption,infection=excluded.infection,
               unusual_activity=excluded.unusual_activity,travel=excluded.travel,medication_change_fact=excluded.medication_change_fact,
               supplement_change_fact=excluded.supplement_change_fact,note=excluded.note,updated_at=excluded.updated_at""",
              (_opaque("check_"),payload["observation_id"],phase_id,payload["day"],payload["adherence"],payload["stress"],payload["sleep_disruption"],payload["infection"],payload["unusual_activity"],payload["travel"],payload["medication_change_fact"],payload["supplement_change_fact"],payload["note"],now,now))
        elif action == "observation_status":
            row=connection.execute("SELECT status FROM personal_observations WHERE id=?",(payload["observation_id"],)).fetchone()
            if row is None:
                raise RuntimeError("status transition not allowed")
            if payload["status"] == str(row["status"]):
                connection.commit()
                return
            if payload["status"] not in TRANSITIONS[str(row["status"])]:
                raise RuntimeError("status transition not allowed")
            if payload["status"] == "active" and row["status"] != "active":
                active_count = int(connection.execute("SELECT count(*) FROM personal_observations WHERE status='active'").fetchone()[0])
                if active_count >= 3:
                    raise RuntimeError("maximum active observations reached")
            connection.execute("UPDATE personal_observations SET status=?,updated_at=? WHERE id=?",(payload["status"],now,payload["observation_id"]))
            if row["status"] != "completed" and payload["status"] == "completed":
                _insert_observation_snapshot(connection, payload["observation_id"], normalized_action_id)
        else:
            row=connection.execute("SELECT status FROM personal_observations WHERE id=?",(payload["observation_id"],)).fetchone()
            if row is None or row["status"] != "completed":
                raise RuntimeError("only completed observations can be recalculated")
            _insert_observation_snapshot(connection,payload["observation_id"],normalized_action_id)
        connection.commit()
    except Exception:
        connection.rollback()
        raise
    finally:
        connection.close()


def apply_document_import(payload: dict[str, Any]) -> str:
    if DASHBOARD_DB is None:
        raise RuntimeError("HEALTH_DASHBOARD_DB is required")
    record, source = load_quarantine(payload["quarantine_token"], DOCUMENT_QUARANTINE)
    if record["sha256"] != payload["sha256"]:
        raise RuntimeError("document hash mismatch")
    metadata = validate_metadata(payload["metadata"])
    if DOCUMENT_STORAGE.exists() and DOCUMENT_STORAGE.is_symlink():
        raise RuntimeError("document storage must not be a symlink")
    DOCUMENT_STORAGE.mkdir(parents=True, mode=0o700, exist_ok=True)
    os.chmod(DOCUMENT_STORAGE, 0o700)
    intake_id = _opaque("doc_")
    now = datetime.now(LOCAL_TIMEZONE).isoformat(timespec="seconds")
    connection = sqlite3.connect(DASHBOARD_DB)
    connection.row_factory = sqlite3.Row
    destination: Path | None = None
    duplicate: sqlite3.Row | None = None
    source_fd: int | None = None
    try:
        connection.execute("PRAGMA foreign_keys=ON")
        apply_document_schema(connection)
        assert_document_schema(connection)
        apply_media_schema(connection)
        assert_media_schema(connection)
        assert_reconciliation_schema(connection)
        connection.commit()
        source_fd = open_quarantine_source(record, source)
        duplicate = connection.execute(
            "SELECT d.id,d.local_original_path FROM dokumente d JOIN document_processing p ON p.document_id=d.id WHERE p.sha256=? ORDER BY d.id LIMIT 1",
            (record["sha256"],),
        ).fetchone()
        if duplicate and duplicate["local_original_path"]:
            destination = Path(str(duplicate["local_original_path"]))
        else:
            destination = DOCUMENT_STORAGE / f"{intake_id}{record['suffix']}"
            assert source_fd is not None
            src_fd = os.dup(source_fd)
            dst_fd = os.open(destination, os.O_WRONLY | os.O_CREAT | os.O_EXCL | os.O_NOFOLLOW, 0o600)
            try:
                with os.fdopen(src_fd, "rb", closefd=False) as src, os.fdopen(dst_fd, "wb", closefd=False) as dst:
                    shutil.copyfileobj(src, dst, length=64 * 1024); dst.flush(); os.fsync(dst.fileno())
            finally:
                os.close(src_fd); os.close(dst_fd)
        connection.execute("BEGIN IMMEDIATE")
        if connection.execute("SELECT 1 FROM sqlite_master WHERE type='table' AND name='dokumente_status'").fetchone():
            connection.execute("INSERT OR IGNORE INTO dokumente_status(status) VALUES('neu')")
        connection.execute(
            """INSERT INTO dokumente(datei_name,dateipfad,daten_typ,groessekbytes,status,kategorie,datei_hash,quelle,review_status,processing_quality,local_original_path,document_date,institution)
               VALUES(?,?,?,?,?,?,?,?,?,?,?,?,?)""",
            (f"{intake_id}{record['suffix']}", str(destination), "pdf" if record["mime"] == "application/pdf" else "image",
             max(1, (int(record["size"]) + 1023) // 1024), "neu", metadata["document_type"], record["sha256"],
             "dashboard_v5_local_intake", "nicht_geprueft", "pending", str(destination), metadata["document_date"] or None,
             metadata["institution"] or None),
        )
        document_id = int(connection.execute("SELECT last_insert_rowid()").fetchone()[0])
        connection.execute(
            """INSERT INTO document_processing(document_id,intake_id,original_status,extraction_status,content_status,search_status,transfer_status,sha256,duplicate_document_id,page_count,investigation_day,personal_title,personal_note,created_at,updated_at)
               VALUES(?,?,'available','not_started','pending','not_searchable','no_candidates',?,?,?,?,?,?,?,?)""",
            (document_id,intake_id,record["sha256"],int(duplicate["id"]) if duplicate else None,record["page_count"],metadata["investigation_day"] or None,metadata["personal_title"] or None,metadata["note"] or None,now,now),
        )
        media_record = record.get("media_info")
        if isinstance(media_record, dict):
            info = MediaInfo(**media_record)
            derivative = {"preview_name":None,"preview_sha256":None,"proxy_name":None,"proxy_sha256":None,"preview_status":"failed"}
            try:
                derivative.update(create_safe_derivatives(destination, DOCUMENT_STORAGE, intake_id, info))
            except (OSError, ValueError, subprocess.TimeoutExpired):
                pass
            connection.execute(
                """INSERT INTO document_media(document_id,media_kind,mime_type,original_sha256,byte_size,width,height,frame_count,auxiliary_count,duration_seconds,container,codec,preview_name,proxy_name,preview_sha256,proxy_sha256,preview_status,decoder_name,decoder_version,probe_version,unconfirmed_captured_at,created_at)
                VALUES(?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)""",
                (document_id,info.media_kind,info.mime_type,info.sha256,info.byte_size,info.width,info.height,info.frame_count,info.auxiliary_count,info.duration_seconds,info.container,info.codec,derivative["preview_name"],derivative["proxy_name"],derivative["preview_sha256"],derivative["proxy_sha256"],derivative["preview_status"],info.decoder_name,info.decoder_version,info.probe_version,info.unconfirmed_captured_at,now),
            )
        connection.commit()
        if duplicate:
            reconcile_documents(connection,[],document_ids=[document_id],ensure_schema=False);connection.commit()
            return intake_id
        try:
            with tempfile.TemporaryDirectory(prefix="health-doc-extract-") as temp:
                engine, engine_version, pages, confidences = extract_pages(destination, record["mime"], Path(temp))
            normalized_pages = [normalize_search_text(page) for page in pages]
            useful = sum(len(page) for page in normalized_pages)
            extraction_status = engine if useful else "no_usable_text"
            if engine == "media_preview":
                connection.execute("BEGIN IMMEDIATE")
                connection.execute("UPDATE dokumente SET extrahierte_inhalte='',processing_quality='media_preview',verarbeite_datum=CURRENT_TIMESTAMP WHERE id=?", (document_id,))
                connection.execute("""UPDATE document_processing SET extraction_status='no_usable_text',content_status='partial',search_status='not_searchable',transfer_status='no_candidates',engine=?,engine_version=?,page_count=0,last_error_code=NULL,updated_at=? WHERE document_id=?""", (engine,engine_version,now,document_id))
                connection.commit()
                reconcile_documents(connection,[],document_ids=[document_id],ensure_schema=False);connection.commit()
                return intake_id
            full_original = "\n\n".join(f"[SEITE {index}]\n{text}" for index, text in enumerate(pages, 1))
            full_normalized = "\n\n".join(f"[SEITE {index}]\n{text}" for index, text in enumerate(normalized_pages, 1))
            version_id = _opaque("text_")
            connection.execute("BEGIN IMMEDIATE")
            connection.execute(
                "INSERT INTO document_text_versions(id,document_id,version,source_kind,original_text,normalized_text,engine,engine_version,created_at) VALUES(?,?,1,?,?,?,?,?,?)",
                (version_id,document_id,engine,full_original,full_normalized,engine,engine_version,now),
            )
            previous = connection.execute(
                """SELECT pg.document_id,pg.original_text,pg.section_hash FROM document_pages pg JOIN dokumente d ON d.id=pg.document_id
                   WHERE pg.document_id<>? AND d.kategorie=? ORDER BY d.document_date DESC,d.id DESC LIMIT 100""",
                (document_id,metadata["document_type"]),
            ).fetchall()
            repeated = 0
            candidate_pages: list[str] = []
            candidate_page_numbers: list[int] = []
            for page_no, (original, normalized, confidence) in enumerate(zip(pages,normalized_pages,confidences),1):
                digest = hashlib.sha256(normalized.casefold().encode()).hexdigest()
                repetition, compared = "first_documented", None
                for old in previous:
                    similarity = section_similarity(original, str(old["original_text"]))
                    if digest == str(old["section_hash"]):
                        repetition, compared, repeated = "identical", int(old["document_id"]), repeated + 1; break
                    if similarity >= .92 and repetition != "identical":
                        repetition, compared = "near_match", int(old["document_id"])
                if previous and repetition == "first_documented": repetition = "new"
                connection.execute(
                    "INSERT INTO document_pages(id,document_id,text_version,page_number,original_text,normalized_text,confidence,section_hash,repetition_status,compared_document_id) VALUES(?,?,?,?,?,?,?,?,?,?)",
                    (_opaque("page_"),document_id,1,page_no,original,normalized,confidence,digest,repetition,compared),
                )
                if repetition != "identical":
                    candidate_pages.append(original)
                    candidate_page_numbers.append(page_no)
            candidates = candidate_rows(candidate_pages, engine, metadata, page_numbers=candidate_page_numbers,identity_scope=intake_id,text_version=1)
            for item in candidates:
                connection.execute(
                    """INSERT OR IGNORE INTO document_candidates(id,document_id,candidate_type,value_text,unit,page_number,section_number,context_text,engine,confidence,source_text_version,candidate_fingerprint,created_at)
                       VALUES(?,?,?,?,?,?,?,?,?,?,?,?,?)""",
                    (item["id"],document_id,item["candidate_type"],item["value_text"],item["unit"],item["page_number"],item["section_number"],item["context_text"],item["engine"],item["confidence"],item["source_text_version"],item["candidate_fingerprint"],now),
                )
            connection.execute("DELETE FROM health_document_machine_fts WHERE document_id=?", (str(document_id),))
            for page_no, normalized in enumerate(normalized_pages, 1):
                if normalized:
                    connection.execute("INSERT INTO health_document_machine_fts(document_id,page_number,section_number,content,review_label) VALUES(?,?,?,?,?)",(str(document_id),page_no,page_no,normalized,"machine_unreviewed"))
            connection.execute("UPDATE dokumente SET extrahierte_inhalte=?,processing_quality=?,verarbeite_datum=CURRENT_TIMESTAMP WHERE id=?",(full_original,extraction_status,document_id))
            connection.execute(
                """UPDATE document_processing SET extraction_status=?,search_status=?,transfer_status=?,engine=?,engine_version=?,page_count=?,updated_at=? WHERE document_id=?""",
                (extraction_status,"machine_searchable" if useful else "not_searchable","candidates_available" if candidates else "no_candidates",engine,engine_version,len(pages),now,document_id),
            )
            connection.commit()
        except Exception:
            connection.rollback()
            connection.execute("UPDATE document_processing SET extraction_status='failed',search_status='not_searchable',last_error_code='extraction_failed',updated_at=? WHERE document_id=?",(now,document_id))
            connection.execute("UPDATE dokumente SET processing_quality='extraction_failed' WHERE id=?",(document_id,))
            connection.commit()
        reconcile_documents(connection,[],document_ids=[document_id],ensure_schema=False)
        connection.commit()
        return intake_id
    except Exception:
        connection.rollback()
        if destination and (not duplicate) and destination.exists(): destination.unlink(missing_ok=True)
        raise
    finally:
        if source_fd is not None:
            os.close(source_fd)
        connection.close()
        source.unlink(missing_ok=True)
        (DOCUMENT_QUARANTINE / f"{payload['quarantine_token']}.json").unlink(missing_ok=True)


def retry_document_extraction(connection: sqlite3.Connection, document_id: int, now: str) -> None:
    document = connection.execute("SELECT local_original_path,dateipfad,document_date,institution FROM dokumente WHERE id=?",(document_id,)).fetchone()
    probe = probe_original(document["local_original_path"] or document["dateipfad"])
    if probe.status != "available" or probe.descriptor is None:
        probe.close(); raise RuntimeError("document original unavailable")
    if str(probe.mime or "").startswith("video/"):
        connection.execute("UPDATE dokumente SET review_status='nicht_geprueft',processing_quality='media_preview' WHERE id=?", (document_id,))
        connection.execute("UPDATE document_processing SET extraction_status='no_usable_text',content_status='partial',search_status='not_searchable',transfer_status='no_candidates',engine='media_preview',last_error_code=NULL,updated_at=? WHERE document_id=?", (now,document_id))
        probe.close()
        return
    try:
        with tempfile.TemporaryDirectory(prefix="health-document-retry-") as directory:
            suffix = ".pdf" if probe.mime == "application/pdf" else ".png" if probe.mime == "image/png" else ".jpg"
            copied = Path(directory) / f"source{suffix}"
            with os.fdopen(os.dup(probe.descriptor),"rb") as source, copied.open("wb") as target:
                shutil.copyfileobj(source,target,64*1024); target.flush(); os.fsync(target.fileno())
            engine, engine_version, pages, confidences = extract_pages(copied, probe.mime or "", Path(directory))
        normalized_pages=[normalize_search_text(page) for page in pages]
        if not any(normalized_pages):
            connection.execute("UPDATE dokumente SET review_status='nicht_geprueft' WHERE id=?",(document_id,))
            connection.execute("UPDATE document_processing SET extraction_status='no_usable_text',content_status='partial',search_status='not_searchable',last_error_code=NULL,updated_at=? WHERE document_id=?",(now,document_id)); return
        version=int(connection.execute("SELECT COALESCE(MAX(version),0)+1 FROM document_text_versions WHERE document_id=?",(document_id,)).fetchone()[0])
        original_text="\n\n".join(f"[SEITE {i}]\n{text}" for i,text in enumerate(pages,1));normalized_text="\n\n".join(f"[SEITE {i}]\n{text}" for i,text in enumerate(normalized_pages,1))
        connection.execute("INSERT INTO document_text_versions(id,document_id,version,source_kind,original_text,normalized_text,engine,engine_version,created_at) VALUES(?,?,?,?,?,?,?,?,?)",(_opaque("text_"),document_id,version,engine,original_text,normalized_text,engine,engine_version,now))
        previous=connection.execute("SELECT document_id,original_text,section_hash FROM document_pages WHERE document_id<>? ORDER BY document_id DESC LIMIT 100",(document_id,)).fetchall()
        intake_scope=str(connection.execute("SELECT intake_id FROM document_processing WHERE document_id=?",(document_id,)).fetchone()[0])
        candidate_pages=[]; candidate_page_numbers=[]
        for page_no,(original,normalized,confidence) in enumerate(zip(pages,normalized_pages,confidences),1):
            repetition="new" if previous else "first_documented"; compared=None
            page_hash=hashlib.sha256(normalized.casefold().encode()).hexdigest()
            for old in previous:
                similarity=section_similarity(original,str(old["original_text"]))
                if page_hash==str(old["section_hash"]): repetition="identical"; compared=int(old["document_id"]); break
                if similarity>=.92: repetition="near_match"; compared=int(old["document_id"])
            connection.execute("INSERT INTO document_pages(id,document_id,text_version,page_number,original_text,normalized_text,confidence,section_hash,repetition_status,compared_document_id) VALUES(?,?,?,?,?,?,?,?,?,?)",(_opaque("page_"),document_id,version,page_no,original,normalized,confidence,page_hash,repetition,compared))
            if repetition != "identical": candidate_pages.append(original); candidate_page_numbers.append(page_no)
        connection.execute("DELETE FROM document_candidates WHERE document_id=? AND status='open'",(document_id,))
        candidates=candidate_rows(candidate_pages,engine,{"document_date":str(document["document_date"] or ""),"institution":str(document["institution"] or "")},page_numbers=candidate_page_numbers,identity_scope=intake_scope,text_version=version)
        for item in candidates:
            connection.execute("INSERT OR IGNORE INTO document_candidates(id,document_id,candidate_type,value_text,unit,page_number,section_number,context_text,engine,confidence,source_text_version,candidate_fingerprint,created_at) VALUES(?,?,?,?,?,?,?,?,?,?,?,?,?)",(item["id"],document_id,item["candidate_type"],item["value_text"],item["unit"],item["page_number"],item["section_number"],item["context_text"],item["engine"],item["confidence"],item["source_text_version"],item["candidate_fingerprint"],now))
        connection.execute("DELETE FROM health_document_machine_fts WHERE document_id=?",(str(document_id),))
        for page_no,normalized in enumerate(normalized_pages,1): connection.execute("INSERT INTO health_document_machine_fts(document_id,page_number,section_number,content,review_label) VALUES(?,?,?,?,?)",(str(document_id),page_no,page_no,normalized,"machine_unreviewed"))
        connection.execute("UPDATE dokumente SET extrahierte_inhalte=?,processing_quality=?,review_status='nicht_geprueft' WHERE id=?",(original_text,engine,document_id))
        connection.execute("UPDATE document_processing SET extraction_status=?,content_status='partial',search_status='machine_searchable',transfer_status=?,page_count=?,engine=?,engine_version=?,last_error_code=NULL,retry_count=retry_count+1,reconciliation_revision=reconciliation_revision+1,updated_at=? WHERE document_id=?",(engine,"candidates_available" if candidates else "no_candidates",len(pages),engine,engine_version,now,document_id))
    finally:
        probe.close()


def _reference_bounds(raw: str | None) -> tuple[str | None,str | None]:
    if not raw:
        return None,None
    number=r"[+-]?\d+(?:[.,]\d+)?"
    text=" ".join(str(raw).split())
    range_match=re.fullmatch(rf"\s*({number})\s*(?:-|–|—|bis)\s*({number})(?:\s+[^\d]*)?\s*",text,re.I)
    if range_match:
        return range_match.group(1).replace(",","."),range_match.group(2).replace(",",".")
    comparator=re.fullmatch(rf"\s*(<=|>=|<|>|≤|≥)\s*({number})(?:\s+[^\d]*)?\s*",text)
    if comparator:
        operator=comparator.group(1).replace("≤","<=").replace("≥",">=")
        value=comparator.group(2).replace(",",".")
        return (None,value) if operator in {"<","<="} else (value,None)
    return None,None


def _reference_unit(raw: str | None) -> str | None:
    if not raw:
        return None
    number=r"[+-]?\d+(?:[.,]\d+)?"
    text=" ".join(str(raw).split())
    match=re.fullmatch(rf"\s*(?:(?:{number})\s*(?:-|–|—|bis)\s*(?:{number})|(?:<=|>=|<|>|≤|≥)\s*(?:{number}))\s*(.*?)\s*",text,re.I)
    suffix=str(match.group(1) if match else "").strip()
    return suffix or None


def _laboratory_value_identity(raw: object) -> str:
    numeric=normalize_value(raw)
    return numeric if numeric is not None else "qual:"+" ".join(str(raw or "").strip().casefold().split())


def _stage_document_candidate(connection: sqlite3.Connection, document_id: int, candidate_id: str, action_id: str, now: str) -> None:
    try:
        preview=transfer_preview(connection,candidate_id)
    except ValueError:
        return
    if preview["target_area"] == "laboratory":
        candidate = connection.execute(
            """SELECT c.value_text,c.context_text,c.unit,m.normalized_parameter,m.normalized_unit
                 FROM document_candidates c JOIN document_candidate_matches m ON m.candidate_id=c.id
                WHERE c.id=? AND c.document_id=?""",
            (candidate_id,document_id),
        ).fetchone()
        observation_date=explicit_candidate_date(candidate["value_text"],candidate["context_text"]) if candidate else None
        if not observation_date:
            raise RuntimeError("laboratory candidate observation date missing")
        preview["new"]["date"]=observation_date
        metric_id,proposed=candidate_catalog_match(candidate["value_text"],candidate["unit"],candidate["normalized_parameter"],candidate["normalized_unit"])
        if not metric_id or not proposed:
            raise RuntimeError("laboratory candidate catalog target missing")
        preview["new"]["parameter"]=proposed
        preview["new"]["unit"]=" ".join(str(candidate["unit"] or "").split()) or None
    stage_id="transfer_"+hashlib.sha256(candidate_id.encode()).hexdigest()[:24]
    connection.execute("""INSERT INTO document_transfer_staging(id,candidate_id,document_id,target_area,status,operation,old_value_json,new_value_json,source_page,expected_text_version,expected_candidate_revision,expected_reconciliation_revision,idempotency_key,reviewed_action_id,created_at)
      VALUES(?,?,?,?,?,?,?,?,?,?,?,?,?,?,?) ON CONFLICT(candidate_id) DO NOTHING""",(stage_id,candidate_id,document_id,preview["target_area"],"reviewed_pending_preview",preview["operation"],_canonical_json(preview["old"]),_canonical_json(preview["new"]),preview["source_page"],preview["expected_text_version"],preview["expected_candidate_revision"],preview["expected_reconciliation_revision"],preview["idempotency_key"],action_id,now))


def _transfer_document_candidate(connection: sqlite3.Connection, document_id: int, candidate_id: str, action_id: str, now: str) -> None:
    stage=connection.execute("SELECT * FROM document_transfer_staging WHERE candidate_id=? AND document_id=?",(candidate_id,document_id)).fetchone()
    if stage is None:raise RuntimeError("candidate not staged")
    if stage["status"]=="transferred":return
    if stage["status"]!="reviewed_pending_preview":raise RuntimeError("candidate transfer blocked")
    current=connection.execute("""SELECT c.status,c.candidate_revision,c.source_text_version,p.reconciliation_revision,p.original_status,p.original_reviewed_at,
      EXISTS(SELECT 1 FROM document_pages pg WHERE pg.document_id=c.document_id
        AND pg.text_version=c.source_text_version AND pg.page_number=c.page_number AND pg.reviewed_at IS NOT NULL) AS source_page_reviewed,
      (SELECT COALESCE(MAX(version),0) FROM document_text_versions WHERE document_id=c.document_id) AS latest_version
      FROM document_candidates c JOIN document_processing p ON p.document_id=c.document_id WHERE c.id=? AND c.document_id=?""",(candidate_id,document_id)).fetchone()
    if current is None or current["status"] not in {"confirmed","corrected_confirmed"}:raise RuntimeError("candidate transfer stale")
    if int(current["candidate_revision"])!=int(stage["expected_candidate_revision"]) or int(current["source_text_version"])!=int(stage["expected_text_version"]) or int(current["latest_version"])!=int(stage["expected_text_version"]) or int(current["reconciliation_revision"])!=int(stage["expected_reconciliation_revision"]):raise RuntimeError("candidate transfer stale")
    if not (current["original_status"]=="available" and current["original_reviewed_at"] and current["source_page_reviewed"]):raise RuntimeError("candidate source review required")
    payload=json.loads(str(stage["new_value_json"]));target=str(stage["target_area"]);canonical_id=None
    if target=="laboratory":
        raw_value=str(payload.get("value") or "")
        parsed=re.search(r"[<>≤≥]?\s*[-+]?\d+(?:[.,]\d+)?",raw_value)
        qualitative=" ".join((raw_value.split(":",1)[1] if ":" in raw_value else raw_value).strip().casefold().split())
        if (not parsed and qualitative not in QUALITATIVE) or not payload.get("date") or not payload.get("parameter"):raise RuntimeError("incomplete laboratory transfer")
        value=parsed.group(0).replace("≤","<=").replace("≥",">=").replace(",",".").replace(" ","") if parsed else qualitative
        public_parameter,public_unit=canonical_lab_pair(payload["parameter"],payload.get("unit"))
        same_parameter=[]
        for existing_row in connection.execute("""SELECT id,parameter_name,wert,einheit,dokument_id,canonical_document_id FROM laborwerte
          WHERE COALESCE(abnahme_datum,befund_datum)=? AND verified_against_original=1""",(str(payload["date"]),)):
            existing_parameter,existing_unit=canonical_lab_pair(existing_row["parameter_name"],existing_row["einheit"])
            if existing_parameter==public_parameter:
                same_parameter.append((int(existing_row["id"]),_laboratory_value_identity(existing_row["wert"]),existing_unit,existing_row["dokument_id"],existing_row["canonical_document_id"]))
        exact=next((item for item in same_parameter if item[1]==_laboratory_value_identity(value) and item[2]==public_unit and document_id in {item[3],item[4]}),None)
        if exact:canonical_id=exact[0]
        elif same_parameter:raise RuntimeError("laboratory transfer conflict")
        else:
            lower,upper=_reference_bounds(payload.get("reference"))
            if payload.get("reference") and lower is None and upper is None:raise RuntimeError("laboratory reference range conflict")
            reference_unit=_reference_unit(payload.get("reference"))
            if reference_unit and normalize_unit(reference_unit)!=normalize_unit(public_unit):raise RuntimeError("laboratory reference unit conflict")
            source_reference=str(payload.get("reference") or "") or None
            cursor=connection.execute("""INSERT INTO laborwerte(dokument_id,parameter_name,wert,einheit,bemerking,reference_min,reference_max,wert_original,quelle,validierungsstatus,source_type,reference_range_source,verified_against_original,canonical_document_id,provenance_note,abnahme_datum,befund_datum)
              VALUES(?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)""",(document_id,str(payload["parameter"]),value,str(payload.get("unit") or ""),f"Quell-Referenzbereich: {source_reference}" if source_reference else None,lower,upper,raw_value,"document_review","validiert","document_candidate","scanned_original",1,document_id,f"reviewed document page {int(stage['source_page'])}; candidate provenance retained",None,str(payload["date"])))
            if cursor.lastrowid is None:raise RuntimeError("laboratory transfer failed")
            canonical_id=int(cursor.lastrowid)
    elif target=="medication":
        if not payload.get("date"):raise RuntimeError("incomplete medication transfer")
        name=str(payload["value"])[:120];dose=str(payload.get("unit") or "")[:40]
        existing=connection.execute("SELECT id FROM medication_administrations WHERE datum=? AND medication_name=? AND COALESCE(dose,'')=? AND event_type='administered' LIMIT 1",(str(payload["date"]),name,dose)).fetchone()
        if existing:canonical_id=int(existing[0])
        else:
            cursor=connection.execute("INSERT INTO medication_administrations(datum,medication_name,dose,route,event_type,notes,source,created_at,occurred_at) VALUES(?,?,?,?,?,?,?,?,?)",(str(payload["date"]),name,dose,"","administered",f"Verified document page {int(stage['source_page'])}","document_review",now,str(payload["date"])+"T12:00:00+02:00"))
            if cursor.lastrowid is None:raise RuntimeError("medication transfer failed")
            canonical_id=int(cursor.lastrowid)
    elif target=="appointment":
        if not payload.get("date"):raise RuntimeError("incomplete appointment transfer")
        institution=connection.execute("SELECT institution FROM dokumente WHERE id=?",(document_id,)).fetchone()[0];reason=str(payload["value"])[:160]
        existing=connection.execute("SELECT id FROM arztbesuche WHERE datum=? AND COALESCE(klinik,'')=? AND grund=? LIMIT 1",(str(payload["date"]),str(institution or "")[:120],reason)).fetchone()
        if existing:canonical_id=int(existing[0])
        else:
            cursor=connection.execute("INSERT INTO arztbesuche(datum,arzt,klinik,grund,zusammenfassung,notizen,created_at) VALUES(?,?,?,?,?,?,?)",(str(payload["date"]),"",str(institution or "")[:120],reason,"",f"Verified document page {int(stage['source_page'])}",now))
            if cursor.lastrowid is None:raise RuntimeError("appointment transfer failed")
            canonical_id=int(cursor.lastrowid)
    else:raise RuntimeError("unsupported transfer target")
    connection.execute("UPDATE document_transfer_staging SET status='transferred',transfer_action_id=?,canonical_row_id=?,transferred_at=? WHERE candidate_id=? AND status='reviewed_pending_preview'",(action_id,canonical_id,now,candidate_id))
    connection.execute("UPDATE document_processing SET transfer_status='partially_transferred',updated_at=? WHERE document_id=?",(now,document_id))


def apply_document_review(payload: dict[str, Any], action_id: str) -> None:
    if DASHBOARD_DB is None:
        raise RuntimeError("HEALTH_DASHBOARD_DB is required")
    connection = sqlite3.connect(DASHBOARD_DB)
    connection.row_factory = sqlite3.Row
    try:
        connection.execute("PRAGMA foreign_keys=ON")
        assert_document_schema(connection)
        apply_media_schema(connection)
        assert_media_schema(connection)
        assert_reconciliation_schema(connection)
        connection.commit()
        connection.execute("BEGIN IMMEDIATE")
        row = connection.execute("SELECT document_id FROM document_processing WHERE intake_id=?",(payload["document_id"],)).fetchone()
        if row is None: raise RuntimeError("document not found")
        document_id = int(row["document_id"])
        digest = hashlib.sha256(_canonical_json(payload).encode()).hexdigest()
        existing = connection.execute("SELECT payload_hash FROM document_review_log WHERE action_id=?",(action_id,)).fetchone()
        if existing:
            if existing["payload_hash"] != digest: raise RuntimeError("action replay mismatch")
            connection.rollback(); return
        operation = payload["operation"]
        now = datetime.now(LOCAL_TIMEZONE).isoformat(timespec="seconds")
        if operation in {"candidate_decision", "candidate_remap", "candidate_transfer"}:
            bound = connection.execute(
                """SELECT c.candidate_type,c.candidate_revision,c.value_text,c.context_text,c.unit,c.source_text_version,c.page_number,
                          m.normalized_parameter,m.normalized_unit,m.reference_text,m.comparison_digest,
                          p.original_status,p.original_reviewed_at,
                          (SELECT pg.reviewed_at FROM document_pages pg
                            WHERE pg.document_id=c.document_id AND pg.text_version=c.source_text_version
                              AND pg.page_number=c.page_number AND pg.reviewed_at IS NOT NULL LIMIT 1) AS source_page_reviewed_at
                     FROM document_candidates c
                     JOIN document_candidate_matches m ON m.candidate_id=c.id
                     JOIN document_processing p ON p.document_id=c.document_id
                    WHERE c.id=? AND c.document_id=?
                      AND c.source_text_version=(SELECT MAX(version) FROM document_text_versions WHERE document_id=?)""",
                (payload["target_id"], document_id, document_id),
            ).fetchone()
            if bound is None:
                raise RuntimeError("candidate revision conflict")
            expected_revision = payload.get("expected_candidate_revision")
            supplied_preview = payload.get("preview_revision")
            requires_binding = str(bound["candidate_type"]) == "laboratory_value"
            if requires_binding and (expected_revision is None or supplied_preview is None):
                raise RuntimeError("candidate preview binding required")
            if requires_binding and operation == "candidate_decision" and payload.get("decision") == "corrected_confirmed":
                raise RuntimeError("laboratory correction requires a regenerated preview")
            if expected_revision is not None and int(bound["candidate_revision"]) != int(expected_revision):
                raise RuntimeError("candidate revision conflict")
            if supplied_preview is not None:
                value_text = " ".join(str(bound["value_text"] or "").split())
                raw_value = value_text.split(":", 1)[1].strip()[:120] if ":" in value_text else ""
                metric_id, proposed = candidate_catalog_match(
                    value_text, bound["unit"], bound["normalized_parameter"], bound["normalized_unit"]
                )
                current_preview = preview_revision_digest(
                    candidate_id=str(payload["target_id"]), candidate_revision=int(bound["candidate_revision"]),
                    comparison_digest=str(bound["comparison_digest"] or ""), target_parameter=proposed,
                    metric_id=metric_id, raw_value=raw_value,
                    unit=" ".join(str(bound["unit"] or bound["normalized_unit"] or "").split()) or None,
                    observation_date=explicit_candidate_date(bound["value_text"], bound["context_text"]),
                    reference_range=" ".join(str(bound["reference_text"] or "").split())[:120] or None,
                    source_original_revision=str(bound["original_reviewed_at"] or "") or None,
                    source_page_revision=str(bound["source_page_reviewed_at"] or "") or None,
                )
                if not hmac.compare_digest(current_preview, str(supplied_preview)):
                    raise RuntimeError("candidate preview revision conflict")
            requires_source_review = operation == "candidate_transfer" or (
                operation == "candidate_decision" and payload.get("decision") == "confirmed"
            )
            if requires_source_review and not (
                bound["original_status"] == "available"
                and bound["original_reviewed_at"]
                and bound["source_page_reviewed_at"]
            ):
                raise RuntimeError("candidate source review required")
        if operation == "retry_extraction":
            retry_document_extraction(connection, document_id, now)
            reconcile_documents(connection,[],document_ids=[document_id],ensure_schema=False)
        elif operation == "original_review":
            changed = connection.execute("UPDATE document_processing SET original_reviewed_at=?,content_status=CASE WHEN content_status='pending' THEN 'partial' ELSE content_status END,updated_at=? WHERE document_id=? AND original_status='available'",(now,now,document_id)).rowcount
            if changed != 1: raise RuntimeError("original unavailable")
        elif operation == "page_review":
            page_number = int(payload["target_id"].split("_",1)[1])
            changed = connection.execute("UPDATE document_pages SET reviewed_at=?,review_action_id=? WHERE document_id=? AND text_version=(SELECT MAX(version) FROM document_text_versions WHERE document_id=?) AND page_number=?",(now,action_id,document_id,document_id,page_number)).rowcount
            if changed != 1: raise RuntimeError("unknown document page")
            connection.execute("UPDATE document_processing SET content_status=CASE WHEN content_status='pending' THEN 'partial' ELSE content_status END,updated_at=? WHERE document_id=?",(now,document_id))
        elif operation == "metadata_review":
            metadata = payload["metadata"]
            connection.execute("UPDATE dokumente SET document_date=?,kategorie=?,institution=? WHERE id=?",(metadata["document_date"] or None,metadata["document_type"],metadata["institution"] or None,document_id))
            connection.execute("UPDATE document_processing SET investigation_day=?,personal_title=?,personal_note=?,updated_at=? WHERE document_id=?",(metadata["investigation_day"] or None,metadata["personal_title"] or None,metadata["note"] or None,now,document_id))
            reconcile_documents(connection,[],document_ids=[document_id],ensure_schema=False)
        elif operation == "candidate_remap":
            selected = next(
                ((parameter, spec) for parameter, spec in LAB_CATALOG_SPECS.items() if spec.metric_id == payload["value"]),
                None,
            )
            if selected is None:
                raise RuntimeError("unknown laboratory target")
            public_parameter, spec = selected
            candidate = connection.execute(
                """SELECT c.value_text,c.unit,c.status,m.normalized_parameter,m.normalized_date,
                          m.normalized_value,m.reference_text
                     FROM document_candidates c JOIN document_candidate_matches m ON m.candidate_id=c.id
                    WHERE c.id=? AND c.document_id=? AND c.candidate_type='laboratory_value'
                      AND c.source_text_version=(SELECT MAX(version) FROM document_text_versions WHERE document_id=?)""",
                (payload["target_id"], document_id, document_id),
            ).fetchone()
            if candidate is None or candidate["status"] not in {"open", "conflicting"}:
                raise RuntimeError("candidate not open")
            if normalize_unit(candidate["unit"]) != normalize_unit(spec.unit):
                raise RuntimeError("laboratory remap unit mismatch")
            normalized_parameter = normalize_parameter(public_parameter)
            normalized_unit = normalize_unit(spec.unit)
            same_parameter = []
            if candidate["normalized_date"]:
                for existing in connection.execute(
                    """SELECT parameter_name,wert,einheit FROM laborwerte
                         WHERE COALESCE(abnahme_datum,befund_datum)=?
                           AND verified_against_original=1""",
                    (candidate["normalized_date"],),
                ):
                    existing_parameter, existing_unit = canonical_lab_pair(existing["parameter_name"], existing["einheit"])
                    if existing_parameter == public_parameter:
                        same_parameter.append((normalize_value(existing["wert"]), normalize_unit(existing_unit), str(existing["wert"])))
            exact = [item for item in same_parameter if item[0] == candidate["normalized_value"] and item[1] == normalized_unit]
            match_status = "exact_match" if exact else "value_conflict" if same_parameter else "not_present"
            next_status = "already_present" if exact else "conflicting" if same_parameter else "open"
            comparison_digest = hashlib.sha256(
                _canonical_json({
                    "candidate": payload["target_id"],
                    "manual_metric": spec.metric_id,
                    "parameter": normalized_parameter,
                    "date": candidate["normalized_date"],
                    "value": candidate["normalized_value"],
                    "unit": normalized_unit,
                    "database_matches": len(same_parameter),
                }).encode()
            ).hexdigest()
            connection.execute(
                """UPDATE document_candidate_matches
                      SET match_status=?,normalized_parameter=?,normalized_unit=?,
                          database_match_count=?,existing_database_value=?,comparison_digest=?,compared_at=?
                    WHERE candidate_id=?""",
                (
                    match_status,
                    normalized_parameter,
                    normalized_unit,
                    len(same_parameter),
                    same_parameter[0][2] if len(same_parameter) == 1 else None,
                    comparison_digest,
                    now,
                    payload["target_id"],
                ),
            )
            connection.execute(
                """UPDATE document_candidates SET status=?,candidate_revision=candidate_revision+1,reviewed_at=?
                    WHERE id=? AND document_id=?""",
                (next_status, now, payload["target_id"], document_id),
            )
            connection.execute(
                """INSERT INTO document_candidate_review_events
                   (action_id,candidate_id,decision,original_value,corrected_value,corrected_unit,processed_at)
                   VALUES(?,?,'corrected',?,?,?,?)""",
                (action_id, payload["target_id"], str(candidate["normalized_parameter"] or ""), public_parameter, spec.unit, now),
            )
            connection.execute(
                "UPDATE document_processing SET reconciliation_revision=reconciliation_revision+1,transfer_status='candidates_available',updated_at=? WHERE document_id=?",
                (now, document_id),
            )
            open_count = int(connection.execute("SELECT COUNT(*) FROM document_candidates WHERE document_id=? AND status IN ('open','conflicting')", (document_id,)).fetchone()[0])
            conflict_count = int(connection.execute("SELECT COUNT(*) FROM document_candidate_matches m JOIN document_candidates c ON c.id=m.candidate_id WHERE c.document_id=? AND m.match_status IN ('value_conflict','ambiguous')", (document_id,)).fetchone()[0])
            connection.execute(
                "UPDATE document_reconciliation SET open_decisions=?,conflict_count=?,queue_bucket='now_reviewable',reason_code=?,prepared_at=? WHERE document_id=?",
                (open_count, conflict_count, "conflict" if conflict_count else "decision_required", now, document_id),
            )
        elif operation == "candidate_decision":
            status_map = {"confirmed":"confirmed","corrected_confirmed":"corrected_confirmed","rejected":"rejected","already_present":"already_present","deferred":"open","conflicting":"conflicting"}
            decision = payload["decision"]
            if decision not in status_map: raise RuntimeError("invalid candidate decision")
            candidate=connection.execute("SELECT value_text,unit,status FROM document_candidates WHERE id=? AND document_id=? AND source_text_version=(SELECT MAX(version) FROM document_text_versions WHERE document_id=?)",(payload["target_id"],document_id,document_id)).fetchone()
            if candidate is None or candidate["status"] not in {"open","conflicting"}:raise RuntimeError("candidate not open")
            event_decision={"confirmed":"correct","corrected_confirmed":"corrected","rejected":"reject","already_present":"already_present","deferred":"defer","conflicting":"defer"}[decision]
            event_value = "__incomplete__" if decision == "conflicting" else payload["value"] if decision == "corrected_confirmed" else None
            connection.execute("INSERT INTO document_candidate_review_events(action_id,candidate_id,decision,original_value,corrected_value,corrected_unit,processed_at) VALUES(?,?,?,?,?,?,?)",(action_id,payload["target_id"],event_decision,str(candidate["value_text"]),event_value,payload["unit"] or None,now))
            if decision!="deferred":
                changed = connection.execute("UPDATE document_candidates SET status=?,corrected_value_text=?,corrected_unit=?,reviewed_at=?,candidate_revision=candidate_revision+1 WHERE id=? AND document_id=? AND source_text_version=(SELECT MAX(version) FROM document_text_versions WHERE document_id=?) AND status IN ('open','conflicting')",(status_map[decision],payload["value"] if decision == "corrected_confirmed" else None,payload["unit"] or None,now,payload["target_id"],document_id,document_id)).rowcount
                if changed != 1: raise RuntimeError("candidate not open")
                if decision in {"confirmed","corrected_confirmed"}:_stage_document_candidate(connection,document_id,payload["target_id"],action_id,now)
            open_count = int(connection.execute("SELECT COUNT(*) FROM document_candidates WHERE document_id=? AND source_text_version=(SELECT MAX(version) FROM document_text_versions WHERE document_id=?) AND status IN ('open','conflicting')",(document_id,document_id)).fetchone()[0])
            staged=connection.execute("SELECT 1 FROM document_transfer_staging WHERE document_id=? AND status IN ('reviewed_pending_preview','ready_for_transfer') LIMIT 1",(document_id,)).fetchone()
            transfer = "candidates_available" if staged or open_count else "rejected"
            connection.execute("UPDATE document_processing SET transfer_status=?,updated_at=? WHERE document_id=?",(transfer,now,document_id))
            reconcile_documents(connection,[],document_ids=[document_id],ensure_schema=False)
        elif operation == "candidate_transfer":
            if payload["decision"]!="transfer":raise RuntimeError("invalid transfer confirmation")
            _transfer_document_candidate(connection,document_id,payload["target_id"],action_id,now)
            connection.execute("INSERT INTO document_candidate_review_events(action_id,candidate_id,decision,original_value,processed_at) SELECT ?,id,'transfer',value_text,? FROM document_candidates WHERE id=? AND document_id=?",(action_id,now,payload["target_id"],document_id))
        elif operation == "content_review":
            readiness = connection.execute("""SELECT p.original_status,p.original_reviewed_at,p.extraction_status,
                (SELECT COUNT(*) FROM document_pages pg WHERE pg.document_id=p.document_id AND pg.text_version=(SELECT MAX(version) FROM document_text_versions WHERE document_id=p.document_id)) AS page_count,
                (SELECT COUNT(*) FROM document_pages pg WHERE pg.document_id=p.document_id AND pg.text_version=(SELECT MAX(version) FROM document_text_versions WHERE document_id=p.document_id) AND pg.reviewed_at IS NULL) AS unreviewed_pages,
                (SELECT COUNT(*) FROM document_candidates c WHERE c.document_id=p.document_id AND c.source_text_version=(SELECT MAX(version) FROM document_text_versions WHERE document_id=p.document_id) AND c.status IN ('open','conflicting')) AS unresolved_candidates,
                (SELECT media_kind FROM document_media m WHERE m.document_id=p.document_id) AS media_kind,
                (SELECT preview_status FROM document_media m WHERE m.document_id=p.document_id) AS preview_status
                FROM document_processing p WHERE p.document_id=?""",(document_id,)).fetchone()
            if readiness is None:
                raise RuntimeError("document review incomplete")
            is_video = readiness["media_kind"] == "video"
            common_incomplete = readiness["original_status"] != "available" or not readiness["original_reviewed_at"] or int(readiness["unresolved_candidates"])
            text_incomplete = not is_video and (readiness["extraction_status"] not in {"text_layer","ocr"} or int(readiness["page_count"]) < 1 or int(readiness["unreviewed_pages"]))
            video_incomplete = is_video and readiness["preview_status"] != "ready"
            if common_incomplete or text_incomplete or video_incomplete:
                raise RuntimeError("document review incomplete")
            connection.execute("UPDATE dokumente SET review_status='geprueft' WHERE id=?",(document_id,))
            connection.execute("UPDATE document_processing SET content_status='reviewed',search_status=?,updated_at=? WHERE document_id=?",("not_searchable" if is_video else "reviewed_searchable",now,document_id))
            if not is_video and connection.execute("SELECT 1 FROM sqlite_master WHERE type='table' AND name='health_document_fts'").fetchone():
                document = connection.execute("SELECT kategorie,institution,daten_typ FROM dokumente WHERE id=?",(document_id,)).fetchone()
                connection.execute("DELETE FROM health_document_fts WHERE document_id=?",(str(document_id),))
                connection.execute("DELETE FROM health_document_repeated_fts WHERE document_id=?",(str(document_id),))
                reviewed_pages=connection.execute("SELECT page_number,normalized_text,repetition_status FROM document_pages WHERE document_id=? AND text_version=(SELECT MAX(version) FROM document_text_versions WHERE document_id=?) ORDER BY page_number",(document_id,document_id)).fetchall()
                for page in reviewed_pages:
                    target="health_document_repeated_fts" if page["repetition_status"]=="identical" else "health_document_fts"
                    connection.execute(f"INSERT INTO {target}(document_id,chunk_no,title,category,institution,document_type,content) VALUES(?,?,?,?,?,?,?)",(str(document_id),int(page["page_number"]),f"Dokument · {str(document['kategorie'] or 'Ohne Kategorie')[:80]}",str(document["kategorie"] or "")[:80],str(document["institution"] or "")[:120],str(document["daten_typ"] or "")[:20],str(page["normalized_text"])))
            connection.execute("DELETE FROM health_document_machine_fts WHERE document_id=?",(str(document_id),))
            reconcile_documents(connection,[],document_ids=[document_id],ensure_schema=False)
        elif operation == "discard_document":
            connection.execute("UPDATE document_processing SET content_status='discarded',transfer_status='rejected',search_status='not_searchable',updated_at=? WHERE document_id=?",(now,document_id))
            connection.execute("DELETE FROM health_document_machine_fts WHERE document_id=?",(str(document_id),))
        else:
            page_number = int(payload["target_id"].split("_",1)[1])
            current_pages = connection.execute(
                "SELECT page_number,original_text,confidence,repetition_status,compared_document_id,reviewed_at,review_action_id FROM document_pages WHERE document_id=? AND text_version=(SELECT MAX(version) FROM document_text_versions WHERE document_id=?) ORDER BY page_number",
                (document_id, document_id),
            ).fetchall()
            if not current_pages or page_number not in {int(item["page_number"]) for item in current_pages}:
                raise RuntimeError("unknown document page")
            version = int(connection.execute("SELECT COALESCE(MAX(version),0)+1 FROM document_text_versions WHERE document_id=?",(document_id,)).fetchone()[0])
            original_parts=[]; normalized_parts=[]
            for page in current_pages:
                page_text = payload["value"] if int(page["page_number"]) == page_number else str(page["original_text"])
                normalized = normalize_search_text(page_text)
                original_parts.append(f"[SEITE {page['page_number']}]\n{page_text}"); normalized_parts.append(f"[SEITE {page['page_number']}]\n{normalized}")
                reviewed_at = None if int(page["page_number"]) == page_number else page["reviewed_at"]
                review_action = None if int(page["page_number"]) == page_number else page["review_action_id"]
                connection.execute("INSERT INTO document_pages(id,document_id,text_version,page_number,original_text,normalized_text,confidence,section_hash,repetition_status,compared_document_id,reviewed_at,review_action_id) VALUES(?,?,?,?,?,?,?,?,?,?,?,?)",(_opaque("page_"),document_id,version,int(page["page_number"]),page_text,normalized,page["confidence"],hashlib.sha256(normalized.casefold().encode()).hexdigest(),page["repetition_status"],page["compared_document_id"],reviewed_at,review_action))
            corrected="\n\n".join(original_parts); normalized="\n\n".join(normalized_parts)
            connection.execute("INSERT INTO document_text_versions(id,document_id,version,source_kind,original_text,normalized_text,engine,engine_version,corrected_by_action_id,created_at) VALUES(?,?,?,'manual_correction',?,?,'manual','1',?,?)",(_opaque("text_"),document_id,version,corrected,normalized,action_id,now))
            connection.execute("DELETE FROM health_document_machine_fts WHERE document_id=?",(str(document_id),))
            for page in current_pages:
                page_text = payload["value"] if int(page["page_number"]) == page_number else str(page["original_text"])
                connection.execute("INSERT INTO health_document_machine_fts(document_id,page_number,section_number,content,review_label) VALUES(?,?,?,?,?)",(str(document_id),int(page["page_number"]),int(page["page_number"]),normalize_search_text(page_text),"machine_corrected_unreviewed"))
            connection.execute("UPDATE dokumente SET extrahierte_inhalte=?,processing_quality='manual_correction' WHERE id=?",(corrected,document_id))
            intake_scope=str(connection.execute("SELECT intake_id FROM document_processing WHERE document_id=?",(document_id,)).fetchone()[0])
            document_metadata=connection.execute("SELECT document_date,institution FROM dokumente WHERE id=?",(document_id,)).fetchone()
            revised_pages=[payload["value"] if int(page["page_number"])==page_number else str(page["original_text"]) for page in current_pages]
            revised_candidates=candidate_rows(revised_pages,"manual",{"document_date":str(document_metadata["document_date"] or ""),"institution":str(document_metadata["institution"] or "")},page_numbers=[int(page["page_number"]) for page in current_pages],identity_scope=intake_scope,text_version=version)
            for item in revised_candidates:
                connection.execute("""INSERT INTO document_candidates(id,document_id,candidate_type,value_text,unit,page_number,section_number,context_text,engine,confidence,status,source_text_version,candidate_fingerprint,created_at)
                    VALUES(?,?,?,?,?,?,?,?,?,?,'open',?,?,?)""",(item["id"],document_id,item["candidate_type"],item["value_text"],item["unit"],item["page_number"],item["section_number"],item["context_text"],item["engine"],item["confidence"],version,item["candidate_fingerprint"],now))
            connection.execute("UPDATE document_processing SET content_status='partial',search_status='machine_searchable',reconciliation_revision=reconciliation_revision+1,updated_at=? WHERE document_id=?",(now,document_id))
            reconcile_documents(connection,[],document_ids=[document_id],ensure_schema=False)
        logged_operation="candidate_decision" if operation in {"candidate_transfer","candidate_remap"} else operation
        connection.execute("INSERT INTO document_review_log(action_id,document_id,action_type,target_id,payload_hash,processed_at) VALUES(?,?,?,?,?,?)",(action_id,document_id,logged_operation,payload["target_id"] or None,digest,now))
        connection.commit()
    except Exception:
        connection.rollback(); raise
    finally: connection.close()


def apply_action(
    payload: dict[str, Any] | str,
    action_id: str | dict[str, int],
    notes: str | None = None,
) -> str | None:
    result_id: str | None = None
    if isinstance(payload, str):
        if not isinstance(action_id, dict) or notes is None:
            raise ValueError("invalid legacy symptom action")
        normalized_payload: dict[str, Any] = {
            "action": "symptom_checkin",
            "date": payload,
            "scores": action_id,
            "notes": notes,
        }
        normalized_action_id = "legacy-symptom-action"
    else:
        normalized_payload = payload
        if not isinstance(action_id, str):
            raise ValueError("invalid action id")
        normalized_action_id = action_id
    if normalized_payload["action"] == "document_import":
        result_id = apply_document_import(normalized_payload)
    elif normalized_payload["action"] == "document_review":
        apply_document_review(normalized_payload, normalized_action_id)
    elif normalized_payload["action"] == "symptom_checkin":
        apply_symptom_action(
            normalized_payload["date"],
            normalized_payload["scores"],
            normalized_payload["notes"],
        )
    elif normalized_payload["action"] == "nutrition_mapping":
        apply_mapping_action(normalized_payload, normalized_action_id)
    elif normalized_payload["action"] == "capture_entry":
        result_id = apply_capture_action(normalized_payload)
    elif normalized_payload["action"] in {
        "symptom_event",
        "medication_event",
        "general_event",
    }:
        apply_patient_action(normalized_payload)
    elif normalized_payload["action"] in {"supplement_plan", "supplement_intake"}:
        apply_supplement_action(normalized_payload)
    elif normalized_payload["action"] in {
        "observation_upsert",
        "observation_phase_upsert",
        "observation_checkin",
        "observation_status",
        "observation_result_snapshot",
    }:
        apply_observation_action(normalized_payload, normalized_action_id)
    else:
        raise ValueError("unsupported action")
    if DASHBOARD_V5_FILE is not None:
        v5_command = [
            sys.executable,
            str(DASHBOARD_V5),
            "--db",
            str(DASHBOARD_DB),
            "--output",
            str(DASHBOARD_V5_FILE),
            "--today",
            local_today().isoformat(),
        ]
        if DASHBOARD_V5_PROFILE == V5_PROFILE_HEALTH_RECORD_6E:
            v5_command.extend(
                ["--health-record-6e", "--explorer-6c", "--calendar-day-6d"]
            )
        subprocess.run(
            v5_command,
            check=True,
            timeout=120,
            capture_output=True,
            text=True,
        )
    return result_id


def write_action_receipt(action_id: str, action: str) -> None:
    if not re.fullmatch(r"[0-9a-f]{32}", action_id):
        raise ValueError("invalid receipt action id")
    if action not in {"nutrition_mapping", "symptom_checkin"}:
        raise ValueError("unsupported receipt action")
    receipt_dir = ACTION_INBOX / "receipts"
    if receipt_dir.exists() and receipt_dir.is_symlink():
        raise RuntimeError("receipt directory must not be a symlink")
    receipt_dir.mkdir(mode=0o700, parents=False, exist_ok=True)
    os.chmod(receipt_dir, 0o700)
    target = receipt_dir / f"{action_id}.json"
    temp = receipt_dir / f".{action_id}.{os.getpid()}.tmp"
    payload = {
        "version": 1,
        "action_id": action_id,
        "action": action,
        "status": "processed",
        "processed_at": datetime.now(LOCAL_TIMEZONE).isoformat(timespec="seconds"),
    }
    descriptor = os.open(temp, os.O_WRONLY | os.O_CREAT | os.O_EXCL | os.O_NOFOLLOW, 0o600)
    try:
        with os.fdopen(descriptor, "w", encoding="utf-8") as handle:
            json.dump(payload, handle, ensure_ascii=True, sort_keys=True, separators=(",", ":"))
            handle.write("\n")
            handle.flush()
            os.fsync(handle.fileno())
        os.replace(temp, target)
        os.chmod(target, 0o600)
    finally:
        try:
            temp.unlink()
        except FileNotFoundError:
            pass


def write_capture_receipt(idempotency_key: str, status: str, action_hash: str) -> None:
    if (
        not re.fullmatch(r"[0-9a-f]{32}", idempotency_key)
        or status not in {"processed", "rejected"}
        or not re.fullmatch(r"[0-9a-f]{64}", action_hash)
    ):
        raise ValueError("invalid capture receipt")
    receipt_dir = ACTION_INBOX / "capture-receipts"
    if receipt_dir.exists() and receipt_dir.is_symlink():
        raise RuntimeError("capture receipt directory must not be a symlink")
    receipt_dir.mkdir(mode=0o700, parents=False, exist_ok=True)
    os.chmod(receipt_dir, 0o700)
    target = receipt_dir / f"{idempotency_key}.json"
    temp = receipt_dir / f".{idempotency_key}.{os.getpid()}.tmp"
    record = {
        "version": 2,
        "status": status,
        "action_hash": action_hash,
        "processed_at": datetime.now(LOCAL_TIMEZONE).isoformat(timespec="seconds"),
    }
    descriptor = os.open(
        temp, os.O_WRONLY | os.O_CREAT | os.O_EXCL | os.O_NOFOLLOW, 0o600
    )
    try:
        with os.fdopen(descriptor, "w", encoding="utf-8") as handle:
            json.dump(
                record, handle, ensure_ascii=True, sort_keys=True, separators=(",", ":")
            )
            handle.write("\n")
            handle.flush()
            os.fsync(handle.fileno())
        os.replace(temp, target)
        os.chmod(target, 0o600)
    finally:
        temp.unlink(missing_ok=True)

def main() -> int:
    if DASHBOARD_DB is None:
        print("HEALTH_DASHBOARD_DB is required", file=sys.stderr)
        return 2
    try:
        media_runtime_self_test()
    except Exception:
        print("local media decoder self-test failed", file=sys.stderr)
        return 2
    ensure_private_inbox()
    failures = 0
    for path in sorted(ACTION_INBOX.glob("*.json")):
        payload: dict[str, Any] | None = None
        receipt_action_hash: str | None = None
        try:
            payload = load_action(path)
            if payload.get("action") == "capture_entry":
                payload = validate_capture_payload(payload)
                receipt_action_hash = hashlib.sha256(_canonical_json(payload).encode()).hexdigest()
            if payload.get("action") == "symptom_checkin":
                apply_action(payload["date"], payload["scores"], payload["notes"])
            else:
                apply_action(payload, path.stem)
            if payload.get("action") in {"nutrition_mapping", "symptom_checkin"}:
                write_action_receipt(path.stem, str(payload["action"]))
            if payload.get("action") == "capture_entry" and receipt_action_hash is not None:
                write_capture_receipt(payload["idempotency_key"], "processed", receipt_action_hash)
            path.unlink()
        except Exception:
            failures += 1
            if payload and payload.get("action") == "capture_entry" and receipt_action_hash is not None:
                try:
                    write_capture_receipt(payload["idempotency_key"], "rejected", receipt_action_hash)
                except Exception:
                    pass
            try:
                path.replace(path.with_suffix(".failed"))
            except OSError:
                pass
    return 1 if failures else 0


if __name__ == "__main__":
    raise SystemExit(main())
