From 9b6e71ce9c61bf5ef79ab765efa986c98d232dce Mon Sep 17 00:00:00 2001 From: Sterling Archer Date: Wed, 12 Aug 2026 15:38:08 -0700 Subject: [PATCH] Add local telemetry logging and 'tune' command to propose rules from recorded history --- solixauto/cli.py | 115 ++++++++++++++++ solixauto/engine.py | 52 +++++++ solixauto/paths.py | 2 + solixauto/tune.py | 327 ++++++++++++++++++++++++++++++++++++++++++++ 4 files changed, 496 insertions(+) create mode 100644 solixauto/tune.py diff --git a/solixauto/cli.py b/solixauto/cli.py index b948beb..0009fba 100644 --- a/solixauto/cli.py +++ b/solixauto/cli.py @@ -1252,6 +1252,98 @@ def load_power_profile(reference): fail(f"{path.name}: {type(err).__name__}: {err}") +def pick_power_profile(): + found = [ + path + for path in list_profiles(paths.POWER_PROFILE_DIR) + if path.name != "README.md" + ] + + if not found: + fail( + "no power profiles found. Run: " + + paths.command("new-profile ") + ) + return None + + if len(found) == 1: + return found[0] + + options = [] + for path in found: + try: + profile = PowerProfile(path) + label = f"{path.name} ({len(profile.active_rules())} rule(s))" + except Exception: + label = f"{path.name} (invalid)" + options.append((path, label)) + + print() + print("Which power profile?") + print() + return choose(options, "Choose") + + +def cmd_tune(args): + from . import tune + + if args.profile: + path = paths.resolve_profile(args.profile, "power") + if path is None: + fail(f"no power profile matching {args.profile!r}") + else: + path = pick_power_profile() + + try: + analysis = tune.Analysis(path, since_hours=args.since_hours) + except tune.TuneError as err: + fail(str(err)) + return + + proposal = analysis.propose() + + print() + print(f"{path.name}") + print() + print(analysis.describe(proposal)) + print() + + rules_yaml = analysis.render_rules_yaml(proposal) + + if args.show_yaml: + print("Proposed rules: block") + print() + print(rules_yaml) + + if args.dry_run: + print("Dry run, nothing was changed.") + return + + print( + "This replaces the rules: block in the profile. The safety floor, " + "notifications, and limits are left untouched." + ) + print() + if not confirm(f"Apply these rules to {path.name}?", default=False): + print("Not applied.") + return + + tune.apply_rules(path, rules_yaml) + + try: + reloaded = PowerProfile(path) + except ProfileError as err: + fail(f"the updated profile failed to validate: {err}") + return + + print() + print(f"Applied. {path.name} now has {len(reloaded.active_rules())} rule(s).") + print( + f"Restart the service for this to take effect: " + f"{paths.command('service ' + path.stem)}" + ) + + def cmd_validate(args): profile = load_power_profile(args.profile) problems, notes = validate(profile) @@ -1575,6 +1667,29 @@ def build_parser(): validate_parser.add_argument("profile") validate_parser.set_defaults(func=cmd_validate) + tune_parser = subparsers.add_parser( + "tune", + help="propose new rule thresholds from recorded telemetry history", + ) + tune_parser.add_argument("profile", nargs="?", default=None) + tune_parser.add_argument( + "--since-hours", + type=float, + default=None, + help="only use telemetry from the last N hours (default: all recorded)", + ) + tune_parser.add_argument( + "--show-yaml", + action="store_true", + help="print the proposed rules: block before asking to apply it", + ) + tune_parser.add_argument( + "--dry-run", + action="store_true", + help="show the proposal and exit without asking to apply anything", + ) + tune_parser.set_defaults(func=cmd_tune) + run_parser = subparsers.add_parser("run", help="run a power profile") run_parser.add_argument("profile") run_parser.add_argument( diff --git a/solixauto/engine.py b/solixauto/engine.py index ef28bda..2f21948 100644 --- a/solixauto/engine.py +++ b/solixauto/engine.py @@ -15,6 +15,8 @@ from .shelly import ShellyTarget EVENT_HISTORY = 200 +TELEMETRY_ROTATE_EVERY = 360 +TELEMETRY_MAX_AGE_DAYS = 90 FIELD_WORDS = { @@ -255,6 +257,7 @@ class Engine: self.last_variables = {} self.last_seen_at = None self.stale_notified = False + self.telemetry_writes = 0 def evaluate(self, variables, now): for state in self.states: @@ -706,6 +709,53 @@ class Engine: except Exception: pass + def record_telemetry(self, serial, variables): + values = { + key: value + for key, value in variables.items() + if isinstance(value, (int, float)) and not isinstance(value, bool) + } + if not values: + return + + line = json.dumps({"t": time.time(), "v": values}) + + try: + paths.TELEMETRY_DIR.mkdir(parents=True, exist_ok=True) + destination = paths.TELEMETRY_DIR / f"{serial}.jsonl" + with open(destination, "a", encoding="utf-8") as handle: + handle.write(line + "\n") + except Exception: + return + + self.telemetry_writes += 1 + if self.telemetry_writes % TELEMETRY_ROTATE_EVERY == 0: + self.rotate_telemetry(destination) + + def rotate_telemetry(self, destination, max_age_days=TELEMETRY_MAX_AGE_DAYS): + cutoff = time.time() - (max_age_days * 86400) + try: + lines = destination.read_text(encoding="utf-8").splitlines() + except Exception: + return + + kept = [] + for line in lines: + try: + record = json.loads(line) + except Exception: + continue + if record.get("t", 0) >= cutoff: + kept.append(line) + + if len(kept) != len(lines): + try: + destination.write_text( + "\n".join(kept) + ("\n" if kept else ""), encoding="utf-8" + ) + except Exception: + pass + def save_state(self, desired, reason): record = { "profile": self.profile.name, @@ -851,6 +901,8 @@ class Engine: self.last_variables = dict(variables) self.last_seen_at = stamp() self.publish_live(variables, age) + serial = self.anker_profile.get("identity", {}).get("serial") or "unknown" + self.record_telemetry(serial, variables) missing = sorted( name diff --git a/solixauto/paths.py b/solixauto/paths.py index 7200f23..a2a7e6d 100644 --- a/solixauto/paths.py +++ b/solixauto/paths.py @@ -13,6 +13,7 @@ SHELLY_PROFILE_DIR = DEVICE_PROFILE_DIR / "shelly" POWER_PROFILE_DIR = BASE_DIR / "power-profiles" STATE_DIR = BASE_DIR / "state" LOG_DIR = BASE_DIR / "logs" +TELEMETRY_DIR = STATE_DIR / "telemetry" RUNTIME_STATE = STATE_DIR / "runtime.json" ENGINE_LOG = LOG_DIR / "automation.log" @@ -25,6 +26,7 @@ ALL_DIRS = [ POWER_PROFILE_DIR, STATE_DIR, LOG_DIR, + TELEMETRY_DIR, ] diff --git a/solixauto/tune.py b/solixauto/tune.py new file mode 100644 index 0000000..c96c491 --- /dev/null +++ b/solixauto/tune.py @@ -0,0 +1,327 @@ +import json +import statistics +import time +from pathlib import Path + +from . import paths +from .profiles import load_yaml +from .rules import PowerProfile, ProfileError + +MIN_HOURS_REQUIRED = 24 +MIN_SAMPLES_REQUIRED = 200 + +TOP_UP_PERCENTILE = 15 +STOP_PERCENTILE = 85 +SURPLUS_RELEASE_PERCENTILE = 60 +SURPLUS_FALLBACK_PERCENTILE = 20 + +FLOOR_MARGIN = 10 +MIN_BAND_WIDTH = 15 + + +class TuneError(Exception): + pass + + +def telemetry_path(serial): + return paths.TELEMETRY_DIR / f"{serial}.jsonl" + + +def load_telemetry(serial, since_epoch=None): + path = telemetry_path(serial) + if not path.exists(): + return [] + + records = [] + for line in path.read_text(encoding="utf-8").splitlines(): + line = line.strip() + if not line: + continue + try: + record = json.loads(line) + except Exception: + continue + t = record.get("t") + if not isinstance(t, (int, float)): + continue + if since_epoch is not None and t < since_epoch: + continue + values = record.get("v") + if isinstance(values, dict): + records.append((t, values)) + + records.sort(key=lambda pair: pair[0]) + return records + + +def series(records, field): + return [ + values[field] + for _, values in records + if field in values and isinstance(values[field], (int, float)) + ] + + +def percentile(values, pct): + if not values: + return None + ordered = sorted(values) + if len(ordered) == 1: + return ordered[0] + rank = (pct / 100) * (len(ordered) - 1) + low = int(rank) + high = min(low + 1, len(ordered) - 1) + frac = rank - low + return ordered[low] + (ordered[high] - ordered[low]) * frac + + +def round_to(value, step): + return round(value / step) * step + + +def describe_span(seconds): + hours = seconds / 3600 + if hours >= 48: + return f"{hours / 24:.1f} days" + if hours >= 1: + return f"{hours:.1f} hours" + return f"{seconds / 60:.0f} minutes" + + +class Analysis: + def __init__(self, profile_path, since_hours=None): + try: + self.profile = PowerProfile(profile_path) + except ProfileError as err: + raise TuneError(f"{profile_path}: {err}") from None + + if self.profile.source_path is None: + raise TuneError( + f"source.profile {self.profile.source_reference!r} does not " + "resolve to a saved Anker device profile" + ) + + anker = load_yaml(self.profile.source_path) + self.serial = (anker.get("identity") or {}).get("serial") + if not self.serial: + raise TuneError( + f"{self.profile.source_path} has no identity.serial" + ) + + self.since_hours = since_hours + since_epoch = None + if since_hours is not None: + since_epoch = time.time() - (since_hours * 3600) + + self.records = load_telemetry(self.serial, since_epoch) + + if len(self.records) < MIN_SAMPLES_REQUIRED: + raise TuneError( + f"only {len(self.records)} telemetry sample(s) recorded for " + f"this device. Need at least {MIN_SAMPLES_REQUIRED} to propose " + "anything sensible. Let the automation run longer, then try " + "again." + ) + + span_seconds = self.records[-1][0] - self.records[0][0] + span_hours = span_seconds / 3600 + + if span_hours < MIN_HOURS_REQUIRED: + raise TuneError( + f"recorded telemetry only spans {describe_span(span_seconds)}. " + f"Need at least {MIN_HOURS_REQUIRED}h of history to propose " + "rules that reflect real usage, not a snapshot. Let the " + "automation run longer, then try again." + ) + + self.span_hours = span_hours + self.span_seconds = span_seconds + + self.battery = series(self.records, "battery_soc") + self.surplus = series(self.records, "pv_surplus") + self.pv_total = series(self.records, "pv_total") + self.load = series(self.records, "output_power_total") + + def daylight_surplus(self): + return [ + values.get("pv_surplus") + for _, values in self.records + if values.get("pv_total", 0) and values.get("pv_total", 0) > 0 + and isinstance(values.get("pv_surplus"), (int, float)) + ] + + def floor_bounds(self): + floor = self.profile.battery_floor + if floor is None: + return None, None + return floor.threshold, floor.release + + def propose(self): + if not self.battery: + raise TuneError( + "no battery_soc readings found in the recorded telemetry" + ) + + floor_at, floor_release = self.floor_bounds() + floor_at = floor_at if floor_at is not None else 0 + floor_release = floor_release if floor_release is not None else floor_at + + low_bound = max(floor_release, floor_at + FLOOR_MARGIN) + + top_up = percentile(self.battery, TOP_UP_PERCENTILE) + stop_at = percentile(self.battery, STOP_PERCENTILE) + + top_up = round_to(top_up, 5) + stop_at = round_to(stop_at, 5) + + top_up = max(top_up, low_bound) + if stop_at - top_up < MIN_BAND_WIDTH: + stop_at = top_up + MIN_BAND_WIDTH + stop_at = min(stop_at, 95) + if stop_at <= top_up: + top_up = max(low_bound, stop_at - MIN_BAND_WIDTH) + + battery_low = percentile(self.battery, 10) + battery_high = percentile(self.battery, 90) + + proposal = { + "top_up_at": top_up, + "stop_at": stop_at, + "observed_low": round(battery_low, 1) if battery_low is not None else None, + "observed_high": round(battery_high, 1) if battery_high is not None else None, + "sample_count": len(self.records), + "span_hours": round(self.span_hours, 1), + "solar": None, + } + + positive_surplus = [v for v in self.daylight_surplus() if v > 0] + if len(positive_surplus) >= 30: + release_surplus = percentile(positive_surplus, SURPLUS_RELEASE_PERCENTILE) + fallback_surplus = percentile(positive_surplus, SURPLUS_FALLBACK_PERCENTILE) + + release_surplus = round_to(release_surplus, 25) + fallback_surplus = round_to(fallback_surplus, 25) + + if fallback_surplus >= release_surplus: + fallback_surplus = max(0, release_surplus - 50) + + proposal["solar"] = { + "release_surplus": release_surplus, + "fallback_surplus": fallback_surplus, + "release_battery": max(top_up, round_to(battery_high or top_up, 5) - 10) + if battery_high + else top_up, + "fallback_battery": top_up, + "sample_count": len(positive_surplus), + } + + return proposal + + def describe(self, proposal): + lines = [] + lines.append( + f"Based on {proposal['sample_count']} sample(s) over " + f"{describe_span(self.span_seconds)}." + ) + lines.append( + f"Battery ranged roughly {proposal['observed_low']:g}% to " + f"{proposal['observed_high']:g}% during that time." + ) + lines.append("") + lines.append( + f"top up from grid when low: battery <= {proposal['top_up_at']:g}%" + ) + lines.append( + f" the battery was at or below this level about " + f"{TOP_UP_PERCENTILE}% of the time recorded" + ) + lines.append( + f"stop charging when full enough: battery >= {proposal['stop_at']:g}%" + ) + lines.append( + f" the battery reached this level or higher about " + f"{100 - STOP_PERCENTILE}% of the time recorded" + ) + + if proposal["solar"]: + solar = proposal["solar"] + lines.append("") + lines.append( + f"solar is carrying it, stay off the grid: " + f"surplus > {solar['release_surplus']:g}W and " + f"battery > {solar['release_battery']:g}%" + ) + lines.append( + f" based on {solar['sample_count']} daylight sample(s); solar " + f"surplus exceeded this level in the upper " + f"{100 - SURPLUS_RELEASE_PERCENTILE}% of observed daylight readings" + ) + lines.append( + f"solar cannot keep up, fall back to the grid: " + f"surplus < {solar['fallback_surplus']:g}W and " + f"battery <= {solar['fallback_battery']:g}%" + ) + lines.append( + f" surplus stayed below this level in the lower " + f"{SURPLUS_FALLBACK_PERCENTILE}% of observed daylight readings" + ) + else: + lines.append("") + lines.append( + "not enough daylight solar data to propose solar rules yet " + "(need at least 30 daylight samples with pv_surplus recorded)" + ) + + return "\n".join(lines) + + def render_rules_yaml(self, proposal, dwell_simple="2m", dwell_solar="15m"): + lines = [] + lines.append("rules:") + lines.append(" - name: top up from grid when low") + lines.append(f" when: battery_soc <= {proposal['top_up_at']:g}") + lines.append(f" for: {dwell_simple}") + lines.append(" then: target.on") + lines.append("") + lines.append(" - name: stop charging when full enough") + lines.append(f" when: battery_soc >= {proposal['stop_at']:g}") + lines.append(f" for: {dwell_simple}") + lines.append(" then: target.off") + + if proposal["solar"]: + solar = proposal["solar"] + lines.append("") + lines.append(" - name: solar is carrying it, stay off the grid") + lines.append( + f" when: pv_surplus > {solar['release_surplus']:g} and " + f"battery_soc > {solar['release_battery']:g}" + ) + lines.append(f" for: {dwell_solar}") + lines.append(" then: target.off") + lines.append("") + lines.append(" - name: solar cannot keep up, fall back to the grid") + lines.append( + f" when: pv_surplus < {solar['fallback_surplus']:g} and " + f"battery_soc <= {solar['fallback_battery']:g}" + ) + lines.append(f" for: {dwell_solar}") + lines.append(" then: target.on") + + return "\n".join(lines) + "\n" + + +def apply_rules(profile_path, rules_yaml): + text = Path(profile_path).read_text(encoding="utf-8") + + start = text.find("\nrules:") + if start == -1: + raise TuneError(f"could not find a rules: block in {profile_path}") + start += 1 + + end = text.find("\nlimits:", start) + if end == -1: + end = len(text) + else: + end += 1 + + new_text = text[:start] + rules_yaml.rstrip("\n") + "\n\n" + text[end:] + Path(profile_path).write_text(new_text, encoding="utf-8")