Add local telemetry logging and 'tune' command to propose rules from recorded history
This commit is contained in:
1 parent
6f46489fb1
commit
9b6e71ce9c
4 files changed
+496
No files matched your search
@@ -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 <name>")
|
||||
)
|
||||
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(
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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,
|
||||
]
|
||||
|
||||
|
||||
|
||||
@@ -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")
|
||||
Reference in new issue
Block a user