"""Datenzugriff für Fahrzeugprofil, Fahrten, Tankvorgänge und Batterieverlauf. Drei getrennte Bestände, wie in SPECIFICATION.md §6.1 festgelegt: - Fahrzeugprofil: eine JSON-Datei, alles Fahrzeugspezifische - Fahrten: JSON Lines, eine Zeile je Fahrt - Tankvorgänge: JSON Lines - Batterieverlauf: JSON Lines, ein Eintrag je Tag Dateiformat und Ordnername sind identisch mit der pyscript-Fassung (/config/audi_dashboard/). Eine bestehende Installation läuft damit nach dem Umstieg einfach weiter - es gibt keine Datenmigration, und deshalb auch keinen Weg, bei der Migration etwas zu verlieren. ZWEI DINGE, DIE HIER ANDERS SIND ALS IN DER PYSCRIPT-FASSUNG ------------------------------------------------------------ 1. Ganz normales `open()`. Die dortige Verrenkung `task.executor(io.open, ...)` war eine reine pyscript-Eigenheit (das eingebaute open() existiert dort nicht). Hier läuft jede Dateioperation stattdessen über `hass.async_add_executor_job` - ein Datei-Zugriff hat im Event-Loop nichts verloren, und diese Klasse ist die einzige Stelle, die das kapselt: alle öffentlichen Methoden sind `async` und schalten selbst in den Executor. 2. Eine Sperre um die verändernden Zugriffe. Fast jede Änderung ist ein Lesen-Ändern-Schreiben über den kompletten Bestand (eine Fahrt ändern heißt: alle Fahrten lesen, eine anfassen, alle zurückschreiben). Liefen zwei davon verschränkt, gewänne die zuletzt schreibende und die andere Änderung wäre spurlos weg. In der pyscript-Fassung war das eine offene Flanke; hier kostet die Absicherung eine `asyncio.Lock`. """ from __future__ import annotations import asyncio import json import logging import os import uuid from collections.abc import Callable from typing import Any from homeassistant.core import HomeAssistant _LOGGER = logging.getLogger(__name__) class Ablage: """Alle Dateizugriffe der App an einer Stelle.""" def __init__(self, hass: HomeAssistant, basis: str) -> None: self._hass = hass self.basis = basis self.profil_pfad = os.path.join(basis, "fahrzeugprofil.json") self.fahrten_pfad = os.path.join(basis, "fahrten.jsonl") self.tankvorgaenge_pfad = os.path.join(basis, "tankvorgaenge.jsonl") self.batterieverlauf_pfad = os.path.join(basis, "batteriespannung.jsonl") self.zuordnung_pfad = os.path.join(basis, "entitaeten.json") self.belege_ordner = os.path.join(basis, "belege") self.backup_ordner = os.path.join(basis, "backups") self._sperre = asyncio.Lock() async def _im_executor(self, funktion: Callable, *args: Any) -> Any: return await self._hass.async_add_executor_job(funktion, *args) # ------------------------------------------------------------- Ordner def _ordner_anlegen(self) -> None: os.makedirs(self.basis, exist_ok=True) os.makedirs(self.belege_ordner, exist_ok=True) async def ordner_sicherstellen(self) -> None: await self._im_executor(self._ordner_anlegen) def _vorlage_anlegen(self, vorlage_pfad: str) -> bool: if os.path.exists(self.profil_pfad): return False with open(vorlage_pfad, encoding="utf-8") as datei: inhalt = datei.read() self._text_schreiben(self.profil_pfad, inhalt) return True async def vorlage_anlegen(self, vorlage_pfad: str) -> bool: """Legt bei der allerersten Einrichtung ein Fahrzeugprofil aus der mitgelieferten Vorlage an. True, wenn dabei wirklich etwas entstanden ist. Vorher war das ein eigener Installationsschritt ("die Datei aus data/ nach /config/audi_dashboard/ kopieren"). Wurde er vergessen, startete die App in einen Zustand, den man von einem Fehler nicht unterscheiden konnte: das Panel blieb leer, im Protokoll stand nur eine Zeile über eine fehlende Datei. Ein bereits vorhandenes Profil wird nie angefasst - deshalb die Existenzprüfung im selben Executor-Aufruf wie das Schreiben, nicht davor: sonst läge zwischen Prüfen und Schreiben ein Zeitfenster.""" return await self._im_executor(self._vorlage_anlegen, vorlage_pfad) # ------------------------------------------------------ Rohe Dateiarbeit @staticmethod def _text_lesen(pfad: str) -> str | None: if not os.path.exists(pfad): return None with open(pfad, encoding="utf-8") as datei: return datei.read() @staticmethod def _text_schreiben(pfad: str, text: str) -> None: """Atomar: erst in eine .tmp-Datei, dann umbenennen. os.replace ist auf allen unterstützten Systemen atomar. Damit kann ein Absturz mitten im Schreiben keine halb geschriebene Datei hinterlassen - der alte Stand bleibt vollständig, bis der neue vollständig da ist.""" os.makedirs(os.path.dirname(pfad), exist_ok=True) tmp = pfad + ".tmp" with open(tmp, "w", encoding="utf-8") as datei: datei.write(text) os.replace(tmp, pfad) @staticmethod def _zeilen_aus_text(text: str | None, pfad: str) -> list[dict]: if not text: return [] datensaetze: list[dict] = [] for nummer, zeile in enumerate(text.splitlines(), start=1): zeile = zeile.strip() if not zeile: continue try: datensaetze.append(json.loads(zeile)) except ValueError as fehler: # Eine kaputte Zeile darf nicht den ganzen Bestand mitreißen: # der Rest ist gültig und die App bleibt benutzbar. _LOGGER.error( "%s Zeile %s ist kein gültiges JSON (%s) - Zeile übersprungen", pfad, nummer, fehler, ) return datensaetze @staticmethod def _text_aus_zeilen(datensaetze: list[dict]) -> str: zeilen = [json.dumps(d, ensure_ascii=False) for d in datensaetze] return "\n".join(zeilen) + ("\n" if zeilen else "") # ------------------------------------------------------- Fahrzeugprofil def _profil_lesen(self) -> dict | None: text = self._text_lesen(self.profil_pfad) if text is None: _LOGGER.error( "%s fehlt - bis dahin bleiben alle Funktionen aus, die das Profil brauchen", self.profil_pfad, ) return None try: return json.loads(text) except ValueError as fehler: _LOGGER.error( "%s ist kein gültiges JSON (%s). Letztes Backup aus %s zurückspielen.", self.profil_pfad, fehler, self.backup_ordner, ) return None async def profil_lesen(self) -> dict | None: """Das Fahrzeugprofil, oder None wenn es fehlt bzw. beschädigt ist. Jeder Aufrufer muss den None-Fall abfangen: ohne diese Prüfung reißt eine fehlende Datei jeden Zeittakt und jeden Dienst mit, der das Profil braucht - bei laufenden Zeittriggern also im Minutentakt.""" return await self._im_executor(self._profil_lesen) async def profil_schreiben(self, profil: dict) -> None: async with self._sperre: await self._im_executor( self._text_schreiben, self.profil_pfad, json.dumps(profil, ensure_ascii=False, indent=2), ) # -------------------------------------------------------------- Fahrten async def fahrten_lesen(self) -> list[dict]: text = await self._im_executor(self._text_lesen, self.fahrten_pfad) return self._zeilen_aus_text(text, self.fahrten_pfad) async def fahrten_schreiben(self, fahrten: list[dict]) -> None: await self._im_executor( self._text_schreiben, self.fahrten_pfad, self._text_aus_zeilen(fahrten) ) async def fahrt_anhaengen(self, fahrt: dict) -> None: async with self._sperre: fahrten = await self.fahrten_lesen() fahrten.append(fahrt) await self.fahrten_schreiben(fahrten) async def fahrten_ergaenzen(self, neue: list[dict]) -> None: """Mehrere Fahrten in einem Rutsch, nach Startzeit sortiert. Für den Historienimport: bei einem Jahr Verlauf wären das sonst hunderte einzelne Schreibvorgänge.""" if not neue: return async with self._sperre: alle = await self.fahrten_lesen() alle.extend(neue) alle.sort(key=lambda f: f.get("ts_start") or "") await self.fahrten_schreiben(alle) @staticmethod def _datensatz_aendern( zeilen: list[dict], id_feld: str, id_wert: str, aenderungen: dict, schutz: bool ) -> bool: """Ersetzt ausgewählte Felder eines Datensatzes anhand seiner ID. Manuell geänderte Felder (edited_fields) werden dabei nie überschrieben. schutz=False hebt genau diese Sperre auf - nötig für Eingaben aus der Oberfläche: edited_fields schützt gegen die automatische Ergänzung, nicht gegen den Menschen, der das Feld gerade selbst korrigiert.""" for d in zeilen: if d.get(id_feld) != id_wert: continue geschuetzt = set(d.get("edited_fields", [])) if schutz else set() for feld, wert in aenderungen.items(): if feld not in geschuetzt: d[feld] = wert return True return False async def fahrt_aktualisieren(self, trip_id: str, aenderungen: dict) -> bool: async with self._sperre: fahrten = await self.fahrten_lesen() if not self._datensatz_aendern(fahrten, "trip_id", trip_id, aenderungen, True): return False await self.fahrten_schreiben(fahrten) return True async def fahrt_bearbeiten(self, trip_id: str, aenderungen: dict) -> bool: """Wie fahrt_aktualisieren(), aber für Eingaben aus der Oberfläche: eine von Hand gesetzte Angabe sticht auch dann, wenn dasselbe Feld schon einmal von Hand gesetzt wurde.""" async with self._sperre: fahrten = await self.fahrten_lesen() if not self._datensatz_aendern(fahrten, "trip_id", trip_id, aenderungen, False): return False await self.fahrten_schreiben(fahrten) return True async def fahrt_loeschen(self, trip_id: str) -> bool: async with self._sperre: fahrten = await self.fahrten_lesen() uebrig = [f for f in fahrten if f.get("trip_id") != trip_id] if len(uebrig) == len(fahrten): return False await self.fahrten_schreiben(uebrig) return True # --------------------------------------------------------- Tankvorgänge async def tankvorgaenge_lesen(self) -> list[dict]: text = await self._im_executor(self._text_lesen, self.tankvorgaenge_pfad) return self._zeilen_aus_text(text, self.tankvorgaenge_pfad) async def tankvorgaenge_schreiben(self, tankvorgaenge: list[dict]) -> None: await self._im_executor( self._text_schreiben, self.tankvorgaenge_pfad, self._text_aus_zeilen(tankvorgaenge), ) async def tankvorgang_anhaengen(self, tankvorgang: dict) -> None: async with self._sperre: alle = await self.tankvorgaenge_lesen() alle.append(tankvorgang) await self.tankvorgaenge_schreiben(alle) async def tankvorgaenge_ergaenzen(self, neue: list[dict]) -> None: if not neue: return async with self._sperre: alle = await self.tankvorgaenge_lesen() alle.extend(neue) alle.sort(key=lambda t: t.get("ts") or "") await self.tankvorgaenge_schreiben(alle) async def tankvorgang_aktualisieren(self, tank_id: str, aenderungen: dict) -> bool: async with self._sperre: alle = await self.tankvorgaenge_lesen() if not self._datensatz_aendern(alle, "tank_id", tank_id, aenderungen, True): return False await self.tankvorgaenge_schreiben(alle) return True async def tankvorgang_loeschen(self, tank_id: str) -> bool: async with self._sperre: alle = await self.tankvorgaenge_lesen() uebrig = [t for t in alle if t.get("tank_id") != tank_id] if len(uebrig) == len(alle): return False await self.tankvorgaenge_schreiben(uebrig) return True async def tankvorgang_nach_id(self, tank_id: str) -> dict | None: for t in await self.tankvorgaenge_lesen(): if t.get("tank_id") == tank_id: return t return None async def tankvorgang_nach_receipt_key(self, receipt_key: str) -> dict | None: """Derselbe Beleg (receipt_key, minutengenau) darf keinen zweiten Datensatz erzeugen (§7.7 Regel 3).""" for t in await self.tankvorgaenge_lesen(): if t.get("receipt_key") == receipt_key: return t return None async def distanz_seit_letzter_tankung(self, aktueller_km: float | None) -> float | None: """Gefahrene Distanz seit dem vorherigen Tankvorgang, als Vorschlag für das gleichnamige Formularfeld - frei überschreibbar, genau wie odometer_km selbst. None, wenn kein Kilometerstand oder kein vorheriger Tankvorgang vorliegt (erster Eintrag überhaupt).""" if aktueller_km is None: return None alle = await self.tankvorgaenge_lesen() if not alle: return None letzter = max(alle, key=lambda t: t.get("ts") or "") if letzter.get("odometer_km") is None: return None return round(aktueller_km - letzter["odometer_km"], 1) # ------------------------------------------------------ Batteriespannung async def batterieverlauf_lesen(self) -> list[dict]: """Ein Eintrag pro Tag ({datum, min, min_ts, max, max_ts, min_temp_c}), älteste zuerst. min_ts/max_ts sind die Zeitstempel (ISO, UTC) der jeweiligen Einzelmessung, für die Datum/Uhrzeit-Anzeige beim Antippen des Diagrammpunkts - der Punkt selbst zeigt nur den Minimalwert. min_temp_c ist die Außentemperatur zum Zeitpunkt genau dieser Minimalmessung (None ohne zugeordneten Temperatursensor oder bei älteren, vor dessen Einführung geschriebenen Einträgen).""" text = await self._im_executor(self._text_lesen, self.batterieverlauf_pfad) return self._zeilen_aus_text(text, self.batterieverlauf_pfad) async def batterieverlauf_tageswert_aktualisieren( self, datum: str, ts: str, spannung: float, aussentemp: float | None = None ) -> None: """Trägt eine neue Messung in den Tageseintrag für `datum` ein: legt ihn beim ersten Wert des Tages an, erweitert sonst nur min/max samt dem Zeitstempel der jeweils neuen Extremmessung. min_temp_c wird nur mitgeschrieben, wenn diese Messung auch den Minimalwert setzt - sie gilt für genau diesen Zeitpunkt, nicht für den Tag insgesamt. Sonderfall, ohne den ein erneuter Import nichts nachliefern könnte: meldet ein späterer Aufruf exakt denselben (bereits gespeicherten) Minimalwert erneut, aber mit einer jetzt verfügbaren Temperatur, wo vorher keine lag - z. B. weil NEBENWERT_MAX_ABSTAND_S seither aufgeweitet wurde und die damals zu weit entfernte Temperaturmessung jetzt innerhalb der Grenze liegt -, wird min_temp_c nachträglich gefüllt. `spannung < eintrag["min"]` allein hätte das nie erlaubt: derselbe Wert ist nie *kleiner* als sich selbst, also griff die Aktualisierung nur beim allerersten Mal.""" async with self._sperre: verlauf = await self.batterieverlauf_lesen() for eintrag in verlauf: if eintrag.get("datum") != datum: continue if spannung < eintrag["min"]: eintrag["min"] = spannung eintrag["min_ts"] = ts eintrag["min_temp_c"] = aussentemp elif spannung == eintrag["min"] and eintrag.get("min_temp_c") is None and aussentemp is not None: eintrag["min_temp_c"] = aussentemp if spannung > eintrag["max"]: eintrag["max"] = spannung eintrag["max_ts"] = ts break else: verlauf.append({ "datum": datum, "min": spannung, "min_ts": ts, "max": spannung, "max_ts": ts, "min_temp_c": aussentemp, }) # Nach Datum sortiert schreiben, nicht in Einfügereihenfolge. # Solange nur die Live-Aufzeichnung schrieb, war beides dasselbe # (sie trägt immer den heutigen Tag ein). Der nachträgliche Import # trägt dagegen vergangene Tage ein - ohne diese Zeile stünden sie # hinter den neueren, und das Diagramm im Frontend, das die Datei # in Dateireihenfolge zeichnet, liefe zeitlich rückwärts. verlauf.sort(key=lambda e: e.get("datum") or "") await self._im_executor( self._text_schreiben, self.batterieverlauf_pfad, self._text_aus_zeilen(verlauf), ) async def batterieverlauf_eintrag_loeschen(self, datum: str) -> bool: async with self._sperre: verlauf = await self.batterieverlauf_lesen() uebrig = [e for e in verlauf if e.get("datum") != datum] if len(uebrig) == len(verlauf): return False await self._im_executor( self._text_schreiben, self.batterieverlauf_pfad, self._text_aus_zeilen(uebrig), ) return True # ---------------------------------------------------- Sensor-Zuordnung async def zuordnung_lesen(self) -> dict: """Die im Setup-Menü gespeicherten Zuordnungen, oder {} wenn die Datei fehlt bzw. beschädigt ist. Der Fehlerfall ist nicht Vorsicht um ihrer selbst willen: das Anwenden läuft beim Start VOR der ersten Veröffentlichung. Ohne die Absicherung reißt eine einzige unlesbare Zeile den gesamten Startvorgang mit - das Panel bliebe komplett leer, ohne dass irgendetwas auf die Ursache hindeutet.""" text = await self._im_executor(self._text_lesen, self.zuordnung_pfad) if not text or not text.strip(): return {} try: gelesen = json.loads(text) except ValueError as fehler: _LOGGER.error( "%s ist kein gültiges JSON (%s). Die eingebauten Standardwerte gelten " "weiter; die Zuordnung lässt sich im Setup-Menü neu speichern.", self.zuordnung_pfad, fehler, ) return {} if not isinstance(gelesen, dict): _LOGGER.error("%s enthält kein Objekt - wird ignoriert.", self.zuordnung_pfad) return {} return gelesen async def zuordnung_schreiben(self, mapping: dict) -> None: await self._im_executor( self._text_schreiben, self.zuordnung_pfad, json.dumps(mapping, ensure_ascii=False, indent=2), ) def neue_id(praefix: str) -> str: return f"{praefix}-{uuid.uuid4().hex[:12]}"