diff --git a/packages/helpermodules/measurement_logging/process_log.py b/packages/helpermodules/measurement_logging/process_log.py index 24eab8feff..fed38e7e98 100644 --- a/packages/helpermodules/measurement_logging/process_log.py +++ b/packages/helpermodules/measurement_logging/process_log.py @@ -1,11 +1,15 @@ from enum import Enum +from copy import deepcopy import json import logging from pathlib import Path from typing import Dict, List, Optional, Tuple, Union +from datetime import datetime +from concurrent.futures import ProcessPoolExecutor +from concurrent.futures.process import BrokenProcessPool from helpermodules import timecheck -from helpermodules.measurement_logging.write_log import (LegacySmartHomeLogData, LogType, create_entry, +from helpermodules.measurement_logging.write_log import (LegacySmartHomeLogData, create_entry, get_previous_entry) from helpermodules.messaging import MessageType, pub_system_message from helpermodules.utils.precision_math import decimal_add, decimal_divide, decimal_multiply, decimal_subtract @@ -245,8 +249,8 @@ def _collect_daily_log_data(date: str): log_data = json.load(json_file) if date == timecheck.create_timestamp_YYYYMMDD(): # beim aktuellen Tag den aktuellen Datensatz ergänzen - log_data["entries"].append(create_entry( - LogType.DAILY, LegacySmartHomeLogData(), get_previous_entry(parent_file, log_data))) + log_data["entries"].append(create_entry(LegacySmartHomeLogData(), + get_previous_entry(parent_file, log_data))) else: # bei älteren als letzten Datensatz den des nächsten Tags try: @@ -263,122 +267,111 @@ def _collect_daily_log_data(date: str): def get_monthly_log(date: str): - data = _collect_monthly_log_data(date) - data["entries"] = _process_entries(data["entries"], CalculationType.ENERGY) - data["totals"] = get_totals(data["entries"], False) - data = _analyse_energy_source(data) - return data + if not (len(date) == 6 and date.isdigit()): + log.debug(f"Ungültiges Datum für Monats-Summen: {date}") + return {"entries": [], "names": {}, "colors": {}, "totals": {}} + # Nur Logs ab dem ältesten Tageslog auswerten + # Sonst werden unnötige Totals-Werte gespeichert + oldest_log_day = _oldest_log_day() + if (oldest_log_day is None + or date < oldest_log_day[:6]): # Jahr und Monat + return {"entries": [], "names": {}, "colors": {}, "totals": {}} -def _collect_monthly_log_data(date: str): - try: - with open(f"{_get_data_folder_path()}/monthly_log/{date}.json", "r") as jsonFile: - log_data = json.load(jsonFile) - this_month = timecheck.create_timestamp_YYYYMM() - if date == this_month: - # add last entry of current day, if current month is requested - try: - today = timecheck.create_timestamp_YYYYMMDD() - with open(f"{_get_data_folder_path()}/daily_log/{today}.json", - "r") as todayJsonFile: - today_log_data = json.load(todayJsonFile) - if len(today_log_data["entries"]) > 0: - log_data["entries"].append(today_log_data["entries"][-1]) - except FILE_ERRORS: - pass - else: - # add first entry of next month - try: - next_date = timecheck.get_relative_date_string(date, month_offset=1) - with open(f"{_get_data_folder_path()}/monthly_log/{next_date}.json", - "r") as nextJsonFile: - next_log_data = json.load(nextJsonFile) - log_data["entries"].append(next_log_data["entries"][0]) - except FILE_ERRORS: - pass - except FILE_ERRORS: - log_data = {"entries": [], "names": {}} - return log_data + monthly_entries = [] + monthly_names = {} + monthly_colors = {} + + this_month = timecheck.create_timestamp_YYYYMM() + today = timecheck.create_timestamp_YYYYMMDD() + day = f"{date}01" + + while day.startswith(date): + # Zukunftstage/-monate nicht verarbeiten, um keine leeren daily_totals zu erzeugen. + if day > today: + break + if day < oldest_log_day: + day = timecheck.get_relative_date_string(day, day_offset=1) + continue + + content = load_daily_source_totals_content(day) + if content is None: + # aktuelle Tageswerte nur berechnen, historische Tage zusaetzlich speichern + content = save_daily_source_totals(day, saving=(day != today)) + + if isinstance(content, dict): + daily_totals = content.get("totals") + daily_entry = content.get("entry") + + if isinstance(daily_totals, dict) and isinstance(daily_entry, dict) and len(daily_entry) > 0: + daily_entry = deepcopy(daily_entry) + daily_entry["date"] = day + _apply_source_totals(daily_entry, daily_totals) + + monthly_entries.append(daily_entry) + if isinstance(content.get("names"), dict): + monthly_names.update(content["names"]) + if isinstance(content.get("colors"), dict): + monthly_colors.update(content["colors"]) + + day = timecheck.get_relative_date_string(day, day_offset=1) + + if len(monthly_entries) > 0: + data = {"entries": monthly_entries, "names": monthly_names, "colors": monthly_colors} + data["totals"] = get_totals(data["entries"], False) + data["totals"] = analyse_percentage_totals(data["entries"], data["totals"]) + data = _analyse_energy_source(data) + + # Fallback für ältere Monate + # da wir den Monat jetzt schon berechnet haben, können wir ihn auch direkt speichern + # falls er noch nicht existiert. + filepath = Path(_get_data_folder_path()) / "monthly_totals" / f"{date}_totals.json" + if not filepath.is_file() and date != this_month and date >= oldest_log_day[:6]: + save_monthly_source_totals(date, data, saving=True) + + return data + + # Fallback, wenn keine Daten vorhanden sind + return {"entries": [], "names": {}, "colors": {}, "totals": {}} def get_yearly_log(year: str): - data = _collect_yearly_log_data(year) - data["entries"] = _process_entries(data["entries"], CalculationType.ENERGY) - data["totals"] = get_totals(data["entries"], False) - data = _analyse_energy_source(data) - return data + # Nur Logs ab dem ältesten Tageslog auswerten + # Sonst werden unnötige totals Werte gespeichert + oldest_log_day = _oldest_log_day() + if oldest_log_day is None or year < oldest_log_day[:4]: + return {"entries": [], "names": {}, "colors": {}, "totals": {}} + results = generate_daily_totals_for_year(year) + monthly_entries = [] + monthly_names = {} + monthly_colors = {} -def _collect_yearly_log_data(year: str): - def add_monthly_log(month: str, check_next_month: bool = False) -> None: - monthly_log_path = Path(__file__).resolve().parents[3]/"data"/"monthly_log" - try: - with open(monthly_log_path / f"{month}.json", "r") as jsonFile: - content = json.load(jsonFile) - entries.append(content["entries"][0]) - # add last entry of current file if next file is missing - if check_next_month: - next_month = timecheck.get_relative_date_string(month, month_offset=1) - if not (monthly_log_path / (next_month+".json")).is_file(): - entries.append(content["entries"][-1]) - log.debug(f"Keine Logdatei für Monat {next_month} gefunden, " - f"füge letzten Datensatz von {month} ein: {entries[-1]['date']}") - names.update(content["names"]) - except FILE_ERRORS: - log.debug(f"Kein Log für Monat {month} gefunden.") - - def add_daily_log(day: str) -> None: - try: - with open(f"{_get_data_folder_path()}/daily_log/{day}.json", "r") as dayJsonFile: - day_log_data = json.load(dayJsonFile) - if len(day_log_data["entries"]) > 0: - entries.append(day_log_data["entries"][-1]) - except FILE_ERRORS: - pass - - entries = [] - names = {} - dates = [] - - # we have to find a valid data range - this_year = timecheck.create_timestamp_YYYY() - this_month = timecheck.create_timestamp_YYYYMM() - if year < this_year: - # if the requested year is in the past, just add all possible months - for month in range(1, 13): - dates.append(f"{year}{month:02}") - else: - # add all months until current month - for month in range(1, int(this_month[-2:])+1): - dates.append(f"{year}{month:02}") - # add data for month range - for date in dates: - try: - log.debug(f"add regular month: {date}") - add_monthly_log(date, date != this_month) - except Exception: - log.exception(f"Fehler beim Zusammenstellen der Jahresdaten für Monat {date}") + for monthly_data in results: + if not isinstance(monthly_data, dict): + continue - # now we have to find a valid "next" entry for proper calculation - if year == this_year: # current year - # add todays last entry - this_day = timecheck.create_timestamp_YYYYMMDD() - try: - log.debug(f"add today: {this_day}") - add_daily_log(this_day) - except Exception: - log.exception(f"Fehler beim Zusammenstellen der Jahresdaten für den aktuellen Tag {this_day}") - else: - # no special handling here, just add first entry of next month - next_date = f"{int(year)+1}01" - try: - log.debug(f"add next month: {next_date}") - add_monthly_log(next_date) - except Exception: - log.exception(f"Fehler beim Zusammenstellen der Jahresdaten für Monat {next_date}") + result_entries = monthly_data.get("entries") + result_names = monthly_data.get("names") + result_colors = monthly_data.get("colors") + + if isinstance(result_entries, list) and len(result_entries) > 0: + monthly_entries.extend(result_entries) + if isinstance(result_names, dict): + monthly_names.update(result_names) + if isinstance(result_colors, dict): + monthly_colors.update(result_colors) + + if len(monthly_entries) > 0: + data = {"entries": monthly_entries, "names": monthly_names, "colors": monthly_colors} + data["totals"] = get_totals(data["entries"], False) + data["totals"] = analyse_percentage_totals(data["entries"], data["totals"]) + data = _analyse_energy_source(data) - # return our data - return {"entries": entries, "names": names} + return data + + # Fallback, wenn keine Daten vorhanden sind + return {"entries": [], "names": {}, "colors": {}, "totals": {}} def _analyse_energy_source(data, calc_cp: Optional[str] = None) -> Dict: @@ -597,13 +590,26 @@ def process_entry(entry: dict, next_entry: dict, calculation: CalculationType): new_data = {} if "imported" in entry[type][module].keys() or "exported" in entry[type][module].keys(): def get_current_and_next(value_key: str) -> Tuple[float, float]: - def get_single_value(source: dict, default: int = 0) -> float: + def get_single_value(source: dict) -> Optional[float]: try: - return source[type][module][value_key] + value = source[type][module][value_key] + if isinstance(value, (int, float)): + return float(value) except KeyError: - return default + pass + return None + current_value = get_single_value(entry) - return current_value, get_single_value(next_entry, current_value) + next_value = get_single_value(next_entry) + + # Keep meter deltas neutral if one side is invalid/missing. + if current_value is None and next_value is None: + return 0.0, 0.0 + if current_value is None: + return next_value, next_value + if next_value is None: + return current_value, current_value + return current_value, next_value value_imported, next_value_imported = get_current_and_next("imported") value_exported, next_value_exported = get_current_and_next("exported") if calculation in [CalculationType.POWER, CalculationType.ALL]: @@ -665,3 +671,244 @@ def _calculate_average_power(time_diff: float, current_imported: float = 0, next def _get_data_folder_path() -> str: return str(Path(__file__).resolve().parents[3] / "data") + + +def save_daily_source_totals(date: str, saving: bool = True): + try: + if not (len(date) == 8 and date.isdigit()): + log.debug(f"Ungültiges Datum für Tages-Summen: {date}") + return None + data = _collect_daily_log_data(date) + source_entries = data.get("entries", []) + processed_entries = _process_entries(deepcopy(source_entries), calculation=CalculationType.ENERGY) + totals = get_totals(processed_entries, process_entries=False) + analysed_data = _analyse_energy_source({ + "entries": processed_entries, + "totals": totals, + "names": data.get("names", {}) + }) + totals = analysed_data["totals"] + + daily_entry = {} + source_daily_entry = _get_last_entry_for_period(source_entries, date, "%Y%m%d") + if source_daily_entry is not None: + daily_entry = deepcopy(source_daily_entry) + daily_entry["date"] = date + _apply_source_totals(daily_entry, totals) + + # Erzeugt Ordner daily_totals, falls nicht vorhanden + totals_dir = Path(_get_data_folder_path()) / "daily_totals" + filepath = totals_dir / f"{date}_totals.json" + + content = { + "date": date, + "totals": totals, + "entry": daily_entry, + "names": data.get("names", {}), + "colors": data.get("colors", {}) + } + + if saving: + totals_dir.mkdir(parents=True, exist_ok=True) + with open(str(filepath), "w") as jsonFile: + json.dump(content, jsonFile, ensure_ascii=False, indent=2) + + log.debug(f"Tages-Summen für {date} gespeichert in {filepath}") + return content + + except FILE_ERRORS: + log.exception(f"Fehler beim Speichern der Tages-Summen für {date}") + + +def load_daily_source_totals_content(date: str): + try: + if not (len(date) == 8 and date.isdigit()): + log.debug(f"Ungültiges Datum für Tages-Summen: {date}") + return None + + filepath = f"{_get_data_folder_path()}/daily_totals/{date}_totals.json" + if not Path(filepath).is_file(): + log.debug(f"Keine Tages-Summen-Datei gefunden: {filepath}") + return None + + with open(str(filepath), "r") as jsonFile: + content = json.load(jsonFile) + + log.debug(f"Tages-Summen für {date} geladen aus {filepath}") + return content + + except FILE_ERRORS: + log.exception(f"Fehler beim Laden der Tages-Summen für {date}") + + +def save_monthly_source_totals(date: str, data: Optional[Dict], saving: bool = True): + try: + # Hauptsächlich für Midnight-Handler + # Wenn keine Daten übergeben werden, dann die Monatswerte berechnen + if data is None: + data = get_monthly_log(date) + + totals = data["totals"] + source_entries = data.get("entries", []) + monthly_entry = {} + source_monthly_entry = _get_last_entry_for_period(source_entries, date, "%Y%m") + if source_monthly_entry is not None: + # Nur den letzten Eintrag des Monats nehmen + monthly_entry = deepcopy(source_monthly_entry) + monthly_entry["date"] = date + + # Erzeugt Ordner monthly_totals, falls nicht vorhanden + totals_dir = Path(_get_data_folder_path()) / "monthly_totals" + filepath = totals_dir / f"{date}_totals.json" + content = { + "date": date, + "totals": totals, + "entry": monthly_entry, + "names": data.get("names", {}), + "colors": data.get("colors", {}) + } + + if saving: + totals_dir.mkdir(parents=True, exist_ok=True) + with open(str(filepath), "w") as jsonFile: + json.dump(content, jsonFile, ensure_ascii=False, indent=2) + + log.debug(f"Monats-Summen für {date} gespeichert in {filepath}") + return content + + except FILE_ERRORS: + log.exception(f"Fehler beim Speichern der Monats-Summen für {date}") + + +def load_monthly_source_totals_content(date: str): + try: + filepath = f"{_get_data_folder_path()}/monthly_totals/{date}_totals.json" + if not Path(filepath).is_file(): + log.debug(f"Keine Monats-Summen-Datei gefunden: {filepath}") + return None + + with open(str(filepath), "r") as jsonFile: + content = json.load(jsonFile) + + log.debug(f"Monats-Summen für {date} geladen aus {filepath}") + return content + + except FILE_ERRORS: + log.exception(f"Fehler beim Laden der Monats-Summen für {date}") + + +def _get_last_entry_for_period(entries: List, period: str, period_format: str) -> Optional[Dict]: + # Suche den letzten Eintrag in der Liste, der dem angegebenen Zeitraum entspricht. + for entry in reversed(entries): + if isinstance(entry, dict) and isinstance(entry.get("timestamp"), (int, float)): + entry_period = datetime.fromtimestamp(entry["timestamp"]).strftime(period_format) + if entry_period == period: + return entry + + # Fallback: Falls kein passender Zeitstempel gefunden wird, letzten gueltigen Eintrag verwenden. + for entry in reversed(entries): + if isinstance(entry, dict): + return entry + return None + + +def _apply_source_totals(entry: Dict, daily_totals: Dict): + for section, section_totals in daily_totals.items(): + section_data = entry.get(section) + if not isinstance(section_totals, dict): + continue + + if not isinstance(section_data, dict): + section_data = {} + entry[section] = section_data + + for module, module_totals in section_totals.items(): + module_data = section_data.get(module) + if not isinstance(module_totals, dict): + continue + + if not isinstance(module_data, dict): + module_data = {} + section_data[module] = module_data + + # Alle vorhandenen Summenfelder des Moduls mit den Tages-Summen ueberschreiben. + module_data.update(module_totals) + return entry + + +def _oldest_log_day() -> Optional[str]: + try: + daily_log_dir = Path(_get_data_folder_path()) / "daily_log" + if not daily_log_dir.is_dir(): + return None + + daily_log_files = [p for p in daily_log_dir.glob("*.json") if p.stem.isdigit()] + if not daily_log_files: + return None + + oldest_file = min(daily_log_files, key=lambda f: f.stem) + return oldest_file.stem + except Exception: + log.exception("Fehler beim Ermitteln des ältesten Tageslogs") + return None + + +def generate_daily_totals_for_year(year: str): + if not (len(year) == 4 and year.isdigit()): + log.debug(f"Ungültiges Jahr für Jahres-Summen: {year}") + return [] + + current_year = timecheck.create_timestamp_YYYY() + if year > current_year: + return [] + + months_list = [] + if current_year == year: + # aktuelles Jahr + current_month = timecheck.create_timestamp_YYYYMM()[4:6] + for month in range(1, int(current_month) + 1): + month_str = f"{current_year}{month:02d}" + months_list.append(month_str) + else: + # historisches Jahr + for month in range(1, 13): + month_str = f"{year}{month:02d}" + months_list.append(month_str) + try: + with ProcessPoolExecutor(max_workers=2) as executor: + results = list(executor.map(get_monthly_parallel, months_list)) + + except BrokenProcessPool: + log.exception(f"Beim vorgenerieren der daily totals fürs Jahr {year} " + f"ist ein Worker-Prozess unerwartet gestorben!") + results = [] + + log.debug(f"Tages-Summen für das Jahr {year} wurden berechnet und gespeichert.") + return results + + +def get_monthly_parallel(month: str): + this_month = timecheck.create_timestamp_YYYYMM() + monthly_entries = [] + monthly_names = {} + monthly_colors = {} + + content = load_monthly_source_totals_content(month) + if content is None: + content = save_monthly_source_totals(month, None, saving=(month != this_month)) + + if isinstance(content, dict): + monthly_totals = content.get("totals") + monthly_entry = content.get("entry") + + if isinstance(monthly_totals, dict) and isinstance(monthly_entry, dict) and len(monthly_entry) > 0: + monthly_entry = deepcopy(monthly_entry) + monthly_entry["date"] = month + _apply_source_totals(monthly_entry, monthly_totals) + + monthly_entries.append(monthly_entry) + if isinstance(content.get("names"), dict): + monthly_names.update(content["names"]) + if isinstance(content.get("colors"), dict): + monthly_colors.update(content["colors"]) + return {"entries": monthly_entries, "names": monthly_names, "colors": monthly_colors} diff --git a/packages/helpermodules/measurement_logging/process_log_unit_test.py b/packages/helpermodules/measurement_logging/process_log_unit_test.py index 30b56ec065..58034dbf98 100644 --- a/packages/helpermodules/measurement_logging/process_log_unit_test.py +++ b/packages/helpermodules/measurement_logging/process_log_unit_test.py @@ -4,6 +4,8 @@ from typing import Dict from unittest.mock import Mock, mock_open import pytest +import datetime +import tempfile from helpermodules.measurement_logging.process_log import ( analyse_percentage, @@ -11,6 +13,9 @@ process_entry, get_totals, _collect_daily_log_data, + _get_last_entry_for_period, + _apply_source_totals, + get_monthly_log, calc_energy_imported_by_source, analyse_percentage_totals, CalculationType) @@ -392,3 +397,155 @@ def test_collect_daily_log_data_json_decode_error(monkeypatch): # evaluation expected_result = {"entries": [], "names": {}} assert result == expected_result + + +def test_get_monthly_log_aggregates_days_and_saves_missing_month_totals(monkeypatch): + # setup + month = "202404" + today = "20240402" + + def relative_date_string(date_value, day_offset=0, month_offset=0): + base = json.loads(json.dumps(date_value)) + if month_offset: + dt = datetime.datetime.strptime(base, "%Y%m") + year = dt.year + ((dt.month - 1 + month_offset) // 12) + month_value = ((dt.month - 1 + month_offset) % 12) + 1 + return f"{year:04d}{month_value:02d}" + dt = datetime.datetime.strptime(base, "%Y%m%d") + return (dt + datetime.timedelta(days=day_offset)).strftime("%Y%m%d") + + # daily_totals Mockdaten + day1_content = { + "totals": {"cp": {"all": {"energy_imported": 10}}}, + "entry": {"timestamp": 1, "cp": {"all": {}}}, + "names": {"cp1": "Ladepunkt 1"}, + "colors": {"cp1": "#123456"} + } + day2_content = { + "totals": {"cp": {"all": {"energy_imported": 20}}}, + "entry": {"timestamp": 2, "cp": {"all": {}}}, + "names": {"cp2": "Ladepunkt 2"}, + "colors": {"cp2": "#654321"} + } + + save_daily_mock = Mock(return_value=day2_content) + save_monthly_mock = Mock() + get_totals_mock = Mock(return_value={"mocked": "totals"}) + + monkeypatch.setattr("helpermodules.measurement_logging.process_log._oldest_log_day", Mock(return_value="20240401")) + + monkeypatch.setattr("helpermodules.measurement_logging.process_log._get_data_folder_path", + Mock(return_value=tempfile.mkdtemp(prefix="process_log_test_"))) + monkeypatch.setattr("helpermodules.measurement_logging.process_log.timecheck.create_timestamp_YYYYMM", + Mock(return_value="202406")) + monkeypatch.setattr("helpermodules.measurement_logging.process_log.timecheck.create_timestamp_YYYYMMDD", + Mock(return_value=today)) + monkeypatch.setattr("helpermodules.measurement_logging.process_log.timecheck.get_relative_date_string", + relative_date_string) + + # Nur für den ersten Tag gibt es bereits eine daily_totals-Datei + monkeypatch.setattr( + "helpermodules.measurement_logging.process_log.load_daily_source_totals_content", + lambda day: day1_content if day == "20240401" else None) + monkeypatch.setattr("helpermodules.measurement_logging.process_log.save_daily_source_totals", save_daily_mock) + monkeypatch.setattr("helpermodules.measurement_logging.process_log.get_totals", get_totals_mock) + monkeypatch.setattr( + "helpermodules.measurement_logging.process_log.analyse_percentage_totals", + lambda entries, totals: {"mocked": "analysed_totals", "count": len(entries), "totals": totals}) + monkeypatch.setattr("helpermodules.measurement_logging.process_log._analyse_energy_source", lambda data: data) + monkeypatch.setattr("helpermodules.measurement_logging.process_log.save_monthly_source_totals", save_monthly_mock) + + # execution + result = get_monthly_log(month) + + # evaluation + assert len(result["entries"]) == 2 + assert result["entries"][0]["date"] == "20240401" + assert result["entries"][0]["cp"]["all"]["energy_imported"] == 10 + assert result["entries"][1]["date"] == "20240402" + assert result["entries"][1]["cp"]["all"]["energy_imported"] == 20 + assert result["names"] == {"cp1": "Ladepunkt 1", "cp2": "Ladepunkt 2"} + assert result["colors"] == {"cp1": "#123456", "cp2": "#654321"} + assert result["totals"] == {"mocked": "analysed_totals", "count": 2, "totals": {"mocked": "totals"}} + + # Für den 2. Tag gibt es noch keine daily_totals-Datei + # -> soll auch nicht gespeichert werden, da der 2. Tag der aktuelle Tag ist + save_daily_mock.assert_called_once_with("20240402", saving=False) + + # monthly_totals sollen gespeichert werden + save_monthly_mock.assert_called_once_with(month, result, saving=True) + + +def test_apply_source_totals_updates_existing_and_creates_missing_sections(): + # setup + entry = { + "cp": { + "all": {"energy_imported": 5, "keep": "x"}, + "cp9": {"keep_cp": True} + }, + "meta": {"unchanged": True} + } + daily_totals = { + "cp": { + "all": {"energy_imported": 12, "energy_exported": 3}, + "cp1": {"energy_imported": 7}, + "cp_invalid": 99 + }, + "hc": { + "all": {"energy_imported": 2} + }, + "invalid_section": "ignore_me" + } + + # execution + result = _apply_source_totals(entry, daily_totals) + + # evaluation + # in-place behavior + assert result is entry + + # existing module gets overwritten/extended, unrelated fields stay + assert result["cp"]["all"]["energy_imported"] == 12 + assert result["cp"]["all"]["energy_exported"] == 3 + assert result["cp"]["all"]["keep"] == "x" + + # missing module and section are created + assert result["cp"]["cp1"]["energy_imported"] == 7 + assert result["hc"]["all"]["energy_imported"] == 2 + + # invalid totals are ignored + assert "cp_invalid" not in result["cp"] + assert "invalid_section" not in result + + # unrelated data stays unchanged + assert result["cp"]["cp9"]["keep_cp"] is True + assert result["meta"]["unchanged"] is True + + +@pytest.mark.parametrize( + "entries, period, period_format, expected", + [ + ( + [ + {"timestamp": 1711929600, "value": "april_1"}, + {"timestamp": 1712016000, "value": "april_2"}, + {"timestamp": 1714608000, "value": "may_2"}, + ], + "202404", + "%Y%m", + {"timestamp": 1712016000, "value": "april_2"}, + ), + ( + [], + "202404", + "%Y%m", + None, + ), + ], +) +def test_get_last_entry_for_period(entries, period, period_format, expected): + # execution + result = _get_last_entry_for_period(entries, period, period_format) + + # evaluation + assert result == expected diff --git a/packages/helpermodules/measurement_logging/update_yields.py b/packages/helpermodules/measurement_logging/update_yields.py index 23b1a8a251..675c5a4e11 100644 --- a/packages/helpermodules/measurement_logging/update_yields.py +++ b/packages/helpermodules/measurement_logging/update_yields.py @@ -1,11 +1,15 @@ -import json import logging from pathlib import Path from typing import Dict from control import data from helpermodules import timecheck -from helpermodules.measurement_logging.process_log import get_totals +from helpermodules.measurement_logging.process_log import (get_totals, + load_daily_source_totals_content, + load_monthly_source_totals_content, + save_daily_source_totals, + get_monthly_log, + generate_daily_totals_for_year) log = logging.getLogger(__name__) @@ -17,6 +21,7 @@ def update_daily_yields(entries): totals = get_totals(entries) [update_module_yields(type, totals) for type in ("bat", "counter", "cp", "pv")] data.data.counter_all_data.data.set.daily_yield_home_consumption = totals["hc"]["all"]["energy_imported"] + return totals except Exception: log.exception("Fehler beim Veröffentlichen der Tageserträge.") @@ -40,72 +45,115 @@ def update_module_yields(module: str, totals: Dict) -> None: log.exception(f"Fehler beim Aktualisieren der Tageserträge für Modul {m} vom Typ {module}.") -def update_pv_monthly_yearly_yields(): - """ veröffentlicht die monatlichen und jährlichen Erträge für PV +def update_pv_monthly_yearly_yields(daily_totals: Dict) -> None: """ - _update_pv_monthly_yields() - _update_pv_yearly_yields() + veröffentlicht die monatlichen und jährlichen Erträge für PV + """ + + folder = _get_parent_path()/"data"/"daily_totals" + if not folder.exists(): + # Nur wenn es noch keine Tages-Totals-Folder gibt, + # werden die totals fürs aktuelle Jahr berechnet und gespeichert. + generate_daily_totals_for_year(timecheck.create_timestamp_YYYY()) + + monthly_totals = _get_pv_monthly_yields(daily_totals) + yearly_totals = _get_pv_yearly_yields(monthly_totals) + + pv_all_monthly_yield = 0 + pv_all_yearly_yield = 0 + + for pv_module in data.data.pv_data.values(): + + # Was wurde im Monat/Jahr exportiert + monthly_yield = monthly_totals.get(f"pv{pv_module.num}", {}).get("energy_exported", 0) + yearly_yield = yearly_totals.get(f"pv{pv_module.num}", {}).get("energy_exported", 0) + + data.data.pv_data[f"pv{pv_module.num}"].data.get.monthly_exported = monthly_yield + data.data.pv_data[f"pv{pv_module.num}"].data.get.yearly_exported = yearly_yield + + # Summe über alle Module für pv_all + pv_all_monthly_yield += monthly_yield + pv_all_yearly_yield += yearly_yield + data.data.pv_all_data.data.get.monthly_exported = pv_all_monthly_yield + data.data.pv_all_data.data.get.yearly_exported = pv_all_yearly_yield -def _update_pv_monthly_yields(): - """ veröffentlicht die monatlichen Erträge für PV - für pv_all nicht die Differenz aus den Logs nehmen, sondern die Summe der Module. Wenn im laufenden Monat ein Modul - gelöscht wurde und keins oder eines mit niedrigerem Zählerstand hinzugefügt wird, wird sonst ein negativer Wert - ermittelt. + +def _get_pv_monthly_yields(daily_totals: Dict) -> Dict: """ - try: - pv_all_monthly_yield = 0 - with open(f"data/monthly_log/{timecheck.create_timestamp_YYYYMM()}.json", "r") as f: - monthly_log = json.load(f) - for pv_module in data.data.pv_data.values(): - for entry in monthly_log["entries"]: - if entry["pv"].get(f"pv{pv_module.num}"): - monthly_yield = data.data.pv_data[f"pv{pv_module.num}"].data.get.exported - \ - entry["pv"][f"pv{pv_module.num}"]["exported"] - pv_all_monthly_yield += monthly_yield - data.data.pv_data[f"pv{pv_module.num}"].data.get.monthly_exported = monthly_yield - break - data.data.pv_all_data.data.get.monthly_exported = pv_all_monthly_yield - except FileNotFoundError: - # am Tag der Ersteinrichtung gibt es noch kein Monatslog-File, das wird erst um Mitternacht erstellt. - log.debug("No monthly logfile found for calculation of monthly yield") - except Exception: - log.exception("Fehler beim Veröffentlichen der monatlichen Erträge für PV") + Berechnet den Unterschied zwischen dem Zählerstand vom ersten Tag des aktuellen Monats bis zum aktuellen Tag. + """ + + this_month = timecheck.create_timestamp_YYYYMM() + today = timecheck.create_timestamp_YYYYMMDD() + + pv_totals = {} + + daily_log_path = _get_parent_path()/"data"/"daily_log" + + for logfile in sorted(daily_log_path.glob(f"{this_month}*.json")): + day = logfile.stem + totals = {} + if day == today: + continue # Der aktuelle Tag wird später behandelt + else: + # Lade alle vergangenen Tage des Monats aus dem daily_totals-Logfile, um die Tageserträge zu ermitteln + content = load_daily_source_totals_content(day) + if content is None: + # Fallback/Migration: fehlende Tages-Totals aus den Tageslogs berechnen und speichern. + content = save_daily_source_totals(day, saving=True) + if isinstance(content, dict): + totals = content.get("totals", {}) + # Totals aufsummieren + _add_pv_totals(pv_totals, totals.get("pv", {})) -def _update_pv_yearly_yields(): - """ veröffentlicht die jährlichen Erträge für PV - für pv_all nicht die Differenz aus den Logs nehmen, sondern die Summe der Module. Wenn unterjährig ein Modul - gelöscht wurde und keins oder eines mit niedrigerem Zählerstand hinzugefügt wird, wird sonst ein negativer Wert - ermittelt. + # aktueller Tag ergänzen + _add_pv_totals(pv_totals, daily_totals.get("pv", {})) + + return pv_totals + + +def _get_pv_yearly_yields(current_monthly_totals: Dict) -> Dict: """ - try: - pv_all_yearly_yield = 0 - path_list = list(Path(_get_parent_path()/"data"/"monthly_log").glob(f"{timecheck.create_timestamp_YYYY()}*")) - sorted_path_list = sorted([str(p) for p in path_list]) - for pv_module in data.data.pv_data.values(): - found_pv = False - for path in sorted_path_list: - with open(path, "r") as f: - monthly_log = json.load(f) - for entry in monthly_log["entries"]: - # erster Eintrag mit PV im Jahr, falls WR erst im laufenden Jahr hinzugefügt wurden - if entry["pv"].get(f"pv{pv_module.num}"): - yearly_yield = data.data.pv_data[f"pv{pv_module.num}"].data.get.exported - \ - entry["pv"][f"pv{pv_module.num}"]["exported"] - pv_all_yearly_yield += yearly_yield - data.data.pv_data[f"pv{pv_module.num}"].data.get.yearly_exported = yearly_yield - found_pv = True - break - if found_pv: - break - else: - # am Tag der Ersteinrichtung gibt es noch kein Monatslog-File, das wird erst um Mitternacht erstellt. - log.debug("No monthly logfile found for calculation of yearly yield") - data.data.pv_all_data.data.get.yearly_exported = pv_all_yearly_yield - except Exception: - log.exception("Fehler beim Veröffentlichen der jährlichen Erträge für PV") + Berechnet den Unterschied zwischen dem Zählerstand vom ersten Monat des aktuellen Jahres bis zum aktuellen Monat. + """ + this_year = timecheck.create_timestamp_YYYY() + this_month = timecheck.create_timestamp_YYYYMM() + + pv_totals = {} + + month = f"{this_year}01" + while month < this_month: + totals = {} + content = load_monthly_source_totals_content(month) + if content is None: + # Fallback/Migration: fehlende Monats-Totals berechnen und speichern. + content = get_monthly_log(month) + if isinstance(content, dict): + totals = content.get("totals", {}) + + # Totals aufsummieren + _add_pv_totals(pv_totals, totals.get("pv", {})) + month = timecheck.get_relative_date_string(month, month_offset=1) + + # aktueller Monat ergänzen + _add_pv_totals(pv_totals, current_monthly_totals) + + return pv_totals def _get_parent_path() -> Path: return Path(__file__).resolve().parents[3] + + +def _add_pv_totals(target: Dict, source: Dict) -> None: + for pv_key, values in source.items(): + energy_exported = values.get("energy_exported", 0) + + # Wenn das PV-Modul noch nicht im target ist, initialisiere es mit 0 + if pv_key not in target: + target[pv_key] = {"energy_exported": 0} + + # Addiere die energy_exported Werte für das PV-Modul + target[pv_key]["energy_exported"] += energy_exported diff --git a/packages/helpermodules/measurement_logging/write_log.py b/packages/helpermodules/measurement_logging/write_log.py index fd8d51aa4b..ea86577707 100644 --- a/packages/helpermodules/measurement_logging/write_log.py +++ b/packages/helpermodules/measurement_logging/write_log.py @@ -1,4 +1,3 @@ -from enum import Enum import os import json import logging @@ -97,11 +96,6 @@ # } -class LogType(Enum): - DAILY = "daily" - MONTHLY = "monthly" - - class LegacySmartHomeLogData: def __init__(self) -> None: self.all_received_topics: Dict = {} @@ -133,20 +127,11 @@ def on_message(self, client: MqttClient, userdata, msg: MQTTMessage): self.all_received_topics.update({msg.topic: msg.payload}) -def save_log(log_type: LogType): - """ Parameter - --------- - folder: str - gibt an, ob ein Tages-oder Monats-Log-Eintrag erstellt werden soll. - """ +def save_log(): try: - parent_file = Path(__file__).resolve().parents[3] / "data" / \ - ("daily_log" if log_type == LogType.DAILY else "monthly_log") + parent_file = Path(__file__).resolve().parents[3] / "data" / "daily_log" parent_file.mkdir(mode=0o755, parents=True, exist_ok=True) - if log_type == LogType.DAILY: - file_name = timecheck.create_timestamp_YYYYMMDD() - else: - file_name = timecheck.create_timestamp_YYYYMM() + file_name = timecheck.create_timestamp_YYYYMMDD() filepath = str(parent_file / f"{file_name}.json") try: @@ -162,7 +147,7 @@ def save_log(log_type: LogType): previous_entry = get_previous_entry(parent_file, content) sh_log_data = LegacySmartHomeLogData() - new_entry = create_entry(log_type, sh_log_data, previous_entry) + new_entry = create_entry(sh_log_data, previous_entry) # json-Objekt in Datei einfügen @@ -194,11 +179,8 @@ def get_previous_entry(parent_file: Path, content: Dict) -> Optional[Dict]: return previous_entry -def create_entry(log_type: LogType, sh_log_data: LegacySmartHomeLogData, previous_entry: Optional[Dict]) -> Dict: - if log_type == LogType.DAILY: - date = timecheck.create_timestamp_HH_MM() - else: - date = timecheck.create_timestamp_YYYYMMDD() +def create_entry(sh_log_data: LegacySmartHomeLogData, previous_entry: Optional[Dict]) -> Dict: + date = timecheck.create_timestamp_HH_MM() current_timestamp = int(timecheck.create_timestamp()) try: diff --git a/packages/main.py b/packages/main.py index 3385139ba9..e2a2d6d066 100755 --- a/packages/main.py +++ b/packages/main.py @@ -27,7 +27,7 @@ from helpermodules.changed_values_handler import ChangedValuesContext from helpermodules.mosquitto_dynsec.mosquitto_dynsec import check_roles_at_start from helpermodules.measurement_logging.update_yields import update_daily_yields, update_pv_monthly_yearly_yields -from helpermodules.measurement_logging.write_log import LogType, save_log +from helpermodules.measurement_logging.write_log import save_log from helpermodules.modbusserver import start_modbus_server from helpermodules.pub import Pub from modules import configuration, loadvars, update_soc @@ -37,6 +37,7 @@ from modules.utils import wait_for_module_update_completed from smarthome.smarthome import readmq, smarthome_handler +from helpermodules.measurement_logging.process_log import save_daily_source_totals, save_monthly_source_totals class HandlerAlgorithm: def __init__(self): @@ -180,9 +181,11 @@ def handler5MinAlgorithm(self): """ try: with ChangedValuesContext(loadvars_.event_module_update_completed): - totals = save_log(LogType.DAILY) - update_daily_yields(totals) - update_pv_monthly_yearly_yields() + entries = save_log() + daily_totals = update_daily_yields(entries) + if daily_totals is not None: + update_pv_monthly_yearly_yields(daily_totals) + for cp in data.data.cp_data.values(): calc_energy_costs(cp) data.data.general_data.grid_protection() @@ -230,7 +233,15 @@ def handler5Min(self): @__with_handler_lock(error_threshold=60) def handler_midnight(self): try: - save_log(LogType.MONTHLY) + today = timecheck.create_timestamp_YYYYMMDD() + previous_day = timecheck.get_relative_date_string(today, day_offset=-1) + save_daily_source_totals(previous_day) + + prev_month = timecheck.get_relative_date_string(today, month_offset=-1)[:6] + # Neuer Monat hat angefangen, daher Monats Totals speichern + if today[6:8] == "01": + save_monthly_source_totals(prev_month, None, saving=True) + thread_errors_path = Path(Path(__file__).resolve().parents[1]/"ramdisk"/"thread_errors.log") with thread_errors_path.open("w") as f: f.write("")