2026-08-10 21:37:46 -07:00
|
|
|
import asyncio
|
|
|
|
|
import json
|
|
|
|
|
import sys
|
|
|
|
|
import time
|
|
|
|
|
from collections import deque
|
|
|
|
|
from datetime import datetime
|
|
|
|
|
|
|
|
|
|
import aiohttp
|
|
|
|
|
|
|
|
|
|
from . import paths
|
2026-08-12 10:53:34 -07:00
|
|
|
from .profiles import load_yaml, slugify
|
2026-08-10 21:37:46 -07:00
|
|
|
from .notify import Notifier, render
|
|
|
|
|
from .rules import derived_values, format_duration, validate
|
|
|
|
|
from .shelly import ShellyTarget
|
|
|
|
|
|
|
|
|
|
|
2026-08-12 10:53:34 -07:00
|
|
|
EVENT_HISTORY = 200
|
2026-08-12 15:38:08 -07:00
|
|
|
TELEMETRY_ROTATE_EVERY = 360
|
|
|
|
|
TELEMETRY_MAX_AGE_DAYS = 90
|
2026-08-12 10:53:34 -07:00
|
|
|
|
|
|
|
|
|
2026-08-12 11:59:04 -07:00
|
|
|
FIELD_WORDS = {
|
|
|
|
|
"pv_surplus": "solar surplus",
|
|
|
|
|
"pv_total": "solar output",
|
|
|
|
|
"output_power_total": "the load",
|
|
|
|
|
"ac_input_power": "grid input",
|
|
|
|
|
"temperature": "temperature",
|
|
|
|
|
"usb_total": "USB output",
|
|
|
|
|
"soc_headroom": "room above the floor",
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
FIELD_UNITS = {
|
|
|
|
|
"pv_surplus": "W",
|
|
|
|
|
"pv_total": "W",
|
|
|
|
|
"output_power_total": "W",
|
|
|
|
|
"ac_input_power": "W",
|
|
|
|
|
"usb_total": "W",
|
|
|
|
|
"temperature": " degrees",
|
|
|
|
|
"soc_headroom": " points",
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
COMPARISON_WORDS = {
|
|
|
|
|
"Lt": "below",
|
|
|
|
|
"LtE": "at or below",
|
|
|
|
|
"Gt": "above",
|
|
|
|
|
"GtE": "at or above",
|
|
|
|
|
"Eq": "exactly",
|
|
|
|
|
"NotEq": "not",
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def describe_clause(node):
|
|
|
|
|
import ast as _ast
|
|
|
|
|
|
|
|
|
|
if not isinstance(node.left, _ast.Name):
|
|
|
|
|
return None
|
|
|
|
|
field = node.left.id
|
|
|
|
|
if field == "battery_soc" or not node.ops or len(node.comparators) != 1:
|
|
|
|
|
return None
|
|
|
|
|
|
|
|
|
|
comparator = node.comparators[0]
|
|
|
|
|
if not isinstance(comparator, _ast.Constant):
|
|
|
|
|
return None
|
|
|
|
|
if not isinstance(comparator.value, (int, float)) or isinstance(
|
|
|
|
|
comparator.value, bool
|
|
|
|
|
):
|
|
|
|
|
return None
|
|
|
|
|
|
|
|
|
|
words = FIELD_WORDS.get(field, field)
|
|
|
|
|
direction = COMPARISON_WORDS.get(type(node.ops[0]).__name__, "at")
|
|
|
|
|
unit = FIELD_UNITS.get(field, "")
|
|
|
|
|
return f"{words} is {direction} {comparator.value:g}{unit}"
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def explain_rule(rule, at, tree):
|
|
|
|
|
import ast as _ast
|
|
|
|
|
|
|
|
|
|
state = rule.desired_state()
|
|
|
|
|
action = "Turns the switch on" if state else "Turns the switch off"
|
|
|
|
|
|
|
|
|
|
direction = "at"
|
|
|
|
|
extras = []
|
|
|
|
|
|
|
|
|
|
for node in _ast.walk(tree):
|
|
|
|
|
if not isinstance(node, _ast.Compare):
|
|
|
|
|
continue
|
|
|
|
|
if isinstance(node.left, _ast.Name) and node.left.id == "battery_soc":
|
|
|
|
|
if node.ops:
|
|
|
|
|
direction = COMPARISON_WORDS.get(
|
|
|
|
|
type(node.ops[0]).__name__, "at"
|
|
|
|
|
)
|
|
|
|
|
continue
|
|
|
|
|
clause = describe_clause(node)
|
|
|
|
|
if clause:
|
|
|
|
|
extras.append(clause)
|
|
|
|
|
|
|
|
|
|
sentence = f"{action} when the battery is {direction} {at:g}%"
|
|
|
|
|
if extras:
|
|
|
|
|
sentence += ", and " + " and ".join(extras)
|
|
|
|
|
sentence += f". Waits {format_duration(rule.dwell)} first."
|
|
|
|
|
return sentence
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def battery_thresholds(rule):
|
|
|
|
|
import ast as _ast
|
|
|
|
|
|
|
|
|
|
points = []
|
|
|
|
|
try:
|
|
|
|
|
tree = _ast.parse(rule.when_source, mode="eval")
|
|
|
|
|
except SyntaxError:
|
|
|
|
|
return points
|
|
|
|
|
|
|
|
|
|
other_fields = {
|
|
|
|
|
node.id
|
|
|
|
|
for node in _ast.walk(tree)
|
|
|
|
|
if isinstance(node, _ast.Name) and node.id != "battery_soc"
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
for node in _ast.walk(tree):
|
|
|
|
|
if not isinstance(node, _ast.Compare):
|
|
|
|
|
continue
|
|
|
|
|
if not isinstance(node.left, _ast.Name) or node.left.id != "battery_soc":
|
|
|
|
|
continue
|
|
|
|
|
if len(node.comparators) != 1:
|
|
|
|
|
continue
|
|
|
|
|
comparator = node.comparators[0]
|
|
|
|
|
if not isinstance(comparator, _ast.Constant):
|
|
|
|
|
continue
|
|
|
|
|
if not isinstance(comparator.value, (int, float)):
|
|
|
|
|
continue
|
|
|
|
|
if isinstance(comparator.value, bool):
|
|
|
|
|
continue
|
|
|
|
|
|
|
|
|
|
state = rule.desired_state()
|
|
|
|
|
value = float(comparator.value)
|
|
|
|
|
points.append(
|
|
|
|
|
{
|
|
|
|
|
"at": value,
|
|
|
|
|
"label": rule.name,
|
|
|
|
|
"kind": "on" if state else ("off" if state is False else "rule"),
|
|
|
|
|
"condition": rule.when_source,
|
|
|
|
|
"explain": explain_rule(rule, value, tree),
|
|
|
|
|
"compound": bool(other_fields),
|
|
|
|
|
}
|
|
|
|
|
)
|
|
|
|
|
|
|
|
|
|
return points
|
|
|
|
|
|
|
|
|
|
|
2026-08-10 21:37:46 -07:00
|
|
|
async def interruptible_sleep(seconds):
|
|
|
|
|
remaining = float(seconds or 0)
|
|
|
|
|
while remaining > 0:
|
|
|
|
|
chunk = min(0.5, remaining)
|
|
|
|
|
await asyncio.sleep(chunk)
|
|
|
|
|
remaining -= chunk
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def configure_event_loop():
|
|
|
|
|
if sys.platform.startswith("win"):
|
|
|
|
|
policy = getattr(asyncio, "WindowsSelectorEventLoopPolicy", None)
|
|
|
|
|
if policy is not None:
|
|
|
|
|
asyncio.set_event_loop_policy(policy())
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def stamp():
|
|
|
|
|
return datetime.now().strftime("%Y-%m-%d %H:%M:%S")
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
class Reporter:
|
|
|
|
|
def __init__(self, log_path=None, quiet=False):
|
|
|
|
|
self.quiet = quiet
|
|
|
|
|
self.handle = None
|
|
|
|
|
if log_path:
|
|
|
|
|
log_path.parent.mkdir(parents=True, exist_ok=True)
|
|
|
|
|
self.handle = open(log_path, "a", encoding="utf-8")
|
|
|
|
|
|
|
|
|
|
def __call__(self, message, force=False):
|
|
|
|
|
line = f"[{stamp()}] {message}"
|
|
|
|
|
if not self.quiet or force:
|
|
|
|
|
print(line, flush=True)
|
|
|
|
|
if self.handle:
|
|
|
|
|
self.handle.write(line + "\n")
|
|
|
|
|
self.handle.flush()
|
|
|
|
|
|
|
|
|
|
def close(self):
|
|
|
|
|
if self.handle:
|
|
|
|
|
self.handle.close()
|
|
|
|
|
self.handle = None
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
class RuleState:
|
|
|
|
|
def __init__(self, rule):
|
|
|
|
|
self.rule = rule
|
|
|
|
|
self.satisfied_since = None
|
|
|
|
|
self.last_value = None
|
|
|
|
|
self.error = None
|
|
|
|
|
|
|
|
|
|
def update(self, variables, now):
|
|
|
|
|
try:
|
|
|
|
|
satisfied = self.rule.evaluate(variables)
|
|
|
|
|
self.error = None
|
|
|
|
|
except Exception as err:
|
|
|
|
|
self.error = f"{type(err).__name__}: {err}"
|
|
|
|
|
satisfied = False
|
|
|
|
|
|
|
|
|
|
self.last_value = satisfied
|
|
|
|
|
|
|
|
|
|
if satisfied:
|
|
|
|
|
if self.satisfied_since is None:
|
|
|
|
|
self.satisfied_since = now
|
|
|
|
|
else:
|
|
|
|
|
self.satisfied_since = None
|
|
|
|
|
|
|
|
|
|
return satisfied
|
|
|
|
|
|
|
|
|
|
def held_for(self, now):
|
|
|
|
|
if self.satisfied_since is None:
|
|
|
|
|
return 0.0
|
|
|
|
|
return now - self.satisfied_since
|
|
|
|
|
|
|
|
|
|
def ripe(self, now):
|
|
|
|
|
if not self.last_value:
|
|
|
|
|
return False
|
|
|
|
|
dwell = self.rule.dwell or 0
|
|
|
|
|
return self.held_for(now) >= dwell
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
class Engine:
|
|
|
|
|
def __init__(self, profile, dry_run=False, reporter=None):
|
|
|
|
|
self.profile = profile
|
|
|
|
|
self.dry_run = dry_run
|
|
|
|
|
self.report = reporter or Reporter()
|
|
|
|
|
|
|
|
|
|
from .anker import AnkerSource
|
|
|
|
|
|
|
|
|
|
self.anker_profile = load_yaml(profile.source_path)
|
|
|
|
|
self.source = AnkerSource(profile.source_path)
|
|
|
|
|
self.target = ShellyTarget(profile.target_path, profile.target_channel)
|
|
|
|
|
|
|
|
|
|
self.notifier = Notifier(
|
|
|
|
|
profile.notifications, reporter=self.report, dry_run=dry_run
|
|
|
|
|
)
|
|
|
|
|
|
|
|
|
|
self.states = [RuleState(rule) for rule in profile.active_rules()]
|
|
|
|
|
self.recent_actions = deque()
|
|
|
|
|
self.last_action_at = 0.0
|
|
|
|
|
self.last_commanded = None
|
|
|
|
|
self.stale_reported = False
|
|
|
|
|
self.floor_latched = False
|
|
|
|
|
self.floor_since = None
|
|
|
|
|
self.last_heartbeat = 0.0
|
|
|
|
|
self.heartbeat_every = 300
|
|
|
|
|
self.incomplete_reported = False
|
2026-08-11 05:55:45 -07:00
|
|
|
self.disconnected_since = None
|
|
|
|
|
self.stale_action_state = None
|
|
|
|
|
self.pre_stale_state = None
|
2026-08-11 05:59:02 -07:00
|
|
|
self.last_variables = {}
|
|
|
|
|
self.last_seen_at = None
|
|
|
|
|
self.stale_notified = False
|
2026-08-12 15:38:08 -07:00
|
|
|
self.telemetry_writes = 0
|
2026-08-10 21:37:46 -07:00
|
|
|
|
|
|
|
|
def evaluate(self, variables, now):
|
|
|
|
|
for state in self.states:
|
|
|
|
|
state.update(variables, now)
|
|
|
|
|
|
|
|
|
|
ripe = [state for state in self.states if state.ripe(now)]
|
|
|
|
|
if not ripe:
|
|
|
|
|
return None, []
|
|
|
|
|
|
|
|
|
|
ripe.sort(key=lambda s: (-s.rule.priority, self.states.index(s)))
|
|
|
|
|
return ripe[0], ripe
|
|
|
|
|
|
|
|
|
|
def rate_limited(self, now):
|
|
|
|
|
while self.recent_actions and now - self.recent_actions[0] > 3600:
|
|
|
|
|
self.recent_actions.popleft()
|
|
|
|
|
|
|
|
|
|
if self.last_action_at and (now - self.last_action_at) < (self.profile.min_gap or 0):
|
|
|
|
|
remaining = (self.profile.min_gap or 0) - (now - self.last_action_at)
|
|
|
|
|
return f"min gap, {format_duration(remaining)} remaining"
|
|
|
|
|
|
|
|
|
|
if len(self.recent_actions) >= self.profile.max_per_hour:
|
|
|
|
|
return f"hourly cap of {self.profile.max_per_hour} actions reached"
|
|
|
|
|
|
|
|
|
|
return None
|
|
|
|
|
|
|
|
|
|
async def apply(
|
2026-08-12 10:53:34 -07:00
|
|
|
self,
|
|
|
|
|
session,
|
|
|
|
|
desired,
|
|
|
|
|
reason,
|
|
|
|
|
now,
|
|
|
|
|
rule=None,
|
|
|
|
|
variables=None,
|
|
|
|
|
force=False,
|
|
|
|
|
cause="rule",
|
2026-08-10 21:37:46 -07:00
|
|
|
):
|
|
|
|
|
if self.last_commanded is desired:
|
|
|
|
|
return False
|
|
|
|
|
|
|
|
|
|
blocked = None if force else self.rate_limited(now)
|
|
|
|
|
if blocked:
|
|
|
|
|
self.report(f"suppressed {self._word(desired)} ({reason}): {blocked}")
|
|
|
|
|
return False
|
|
|
|
|
|
|
|
|
|
if self.dry_run:
|
|
|
|
|
self.report(f"DRY RUN would turn {self._word(desired)} - {reason}")
|
|
|
|
|
self.last_commanded = desired
|
|
|
|
|
return True
|
|
|
|
|
|
|
|
|
|
try:
|
|
|
|
|
await self.target.set_state(session, desired)
|
|
|
|
|
except Exception as err:
|
|
|
|
|
self.report(f"FAILED to turn {self._word(desired)}: {type(err).__name__}: {err}")
|
|
|
|
|
return False
|
|
|
|
|
|
|
|
|
|
self.last_commanded = desired
|
|
|
|
|
self.last_action_at = now
|
|
|
|
|
self.recent_actions.append(now)
|
|
|
|
|
self.report(f"turned {self._word(desired)} - {reason}")
|
|
|
|
|
self.save_state(desired, reason)
|
2026-08-12 10:53:34 -07:00
|
|
|
self.record_event(desired, reason, rule, variables, cause=cause)
|
2026-08-10 21:37:46 -07:00
|
|
|
await self.notify(desired, rule, variables, event="action")
|
|
|
|
|
return True
|
|
|
|
|
|
|
|
|
|
@staticmethod
|
|
|
|
|
def _word(desired):
|
|
|
|
|
return "ON" if desired else "OFF"
|
|
|
|
|
|
|
|
|
|
def context(self, desired, rule, variables, event="action"):
|
|
|
|
|
source_identity = self.anker_profile.get("identity", {})
|
|
|
|
|
target_identity = self.target.profile.get("identity", {})
|
|
|
|
|
|
|
|
|
|
context = dict(variables or {})
|
|
|
|
|
context.update(
|
|
|
|
|
{
|
|
|
|
|
"profile": self.profile.name,
|
|
|
|
|
"event": event,
|
|
|
|
|
"rule": rule.name if rule else "",
|
|
|
|
|
"condition": rule.when_source if rule else "",
|
|
|
|
|
"action": self._word(desired) if desired is not None else "",
|
|
|
|
|
"action_word": ("on" if desired else "off") if desired is not None else "",
|
|
|
|
|
"source_name": (
|
|
|
|
|
source_identity.get("name")
|
|
|
|
|
or source_identity.get("model")
|
|
|
|
|
or source_identity.get("serial")
|
|
|
|
|
),
|
|
|
|
|
"source_model": source_identity.get("model", ""),
|
|
|
|
|
"source_serial": source_identity.get("serial", ""),
|
|
|
|
|
"target_name": (
|
|
|
|
|
target_identity.get("name")
|
|
|
|
|
or target_identity.get("model")
|
|
|
|
|
or self.target.host
|
|
|
|
|
),
|
|
|
|
|
"target_model": target_identity.get("model", ""),
|
|
|
|
|
"target_host": self.target.host,
|
|
|
|
|
"target_channel": self.target.channel,
|
|
|
|
|
"time": stamp(),
|
|
|
|
|
}
|
|
|
|
|
)
|
|
|
|
|
return context
|
|
|
|
|
|
2026-08-11 05:59:02 -07:00
|
|
|
async def notify(self, desired, rule, variables, event="action", extra=None):
|
2026-08-10 21:37:46 -07:00
|
|
|
settings = self.profile.notifications
|
2026-08-11 05:59:02 -07:00
|
|
|
subscription = "stale" if event == "recovered" else event
|
|
|
|
|
if not settings.wants(subscription):
|
2026-08-10 21:37:46 -07:00
|
|
|
return
|
|
|
|
|
if event == "action" and not settings.rule_enabled(rule):
|
|
|
|
|
return
|
|
|
|
|
|
2026-08-11 05:59:02 -07:00
|
|
|
if event in ("stale", "recovered"):
|
|
|
|
|
merged = dict(self.last_variables)
|
|
|
|
|
merged.update(variables or {})
|
|
|
|
|
context = self.context(desired, rule, merged, event)
|
|
|
|
|
context["last_seen"] = self.last_seen_at or "unknown"
|
|
|
|
|
context.update(extra or {})
|
|
|
|
|
template = (
|
|
|
|
|
settings.stale_template
|
|
|
|
|
if event == "stale"
|
|
|
|
|
else settings.recovered_template
|
|
|
|
|
)
|
|
|
|
|
body = render(template, context)
|
|
|
|
|
title = render(settings.title, context)
|
|
|
|
|
await self.notifier.send(
|
|
|
|
|
title, body, priority=settings.priority, key=event, force=True
|
|
|
|
|
)
|
|
|
|
|
return
|
|
|
|
|
|
2026-08-10 21:37:46 -07:00
|
|
|
context = self.context(desired, rule, variables, event)
|
|
|
|
|
body = render(settings.template_for(rule), context)
|
|
|
|
|
title = render(settings.title, context)
|
|
|
|
|
priority = (rule.notify_priority if rule else None) or settings.priority
|
|
|
|
|
key = rule.name if rule else event
|
|
|
|
|
|
|
|
|
|
await self.notifier.send(title, body, priority=priority, key=key)
|
|
|
|
|
|
|
|
|
|
def required_fields(self):
|
|
|
|
|
names = set(self.profile.referenced_names())
|
|
|
|
|
floor = self.profile.battery_floor
|
|
|
|
|
if floor is not None and floor.enabled:
|
|
|
|
|
names.add(floor.field)
|
|
|
|
|
return names
|
|
|
|
|
|
|
|
|
|
def summarize(self, variables, now):
|
|
|
|
|
names = sorted(self.profile.referenced_names())
|
|
|
|
|
floor = self.profile.battery_floor
|
|
|
|
|
if floor and floor.field not in names:
|
|
|
|
|
names.insert(0, floor.field)
|
|
|
|
|
|
2026-08-10 22:24:01 -07:00
|
|
|
parts_readings = [
|
2026-08-10 21:37:46 -07:00
|
|
|
f"{name}={variables.get(name)}" for name in names if name in variables
|
2026-08-10 22:24:01 -07:00
|
|
|
]
|
|
|
|
|
|
|
|
|
|
for name in self.profile.monitor_fields:
|
|
|
|
|
if name in names:
|
|
|
|
|
continue
|
|
|
|
|
value = variables.get(name, None)
|
|
|
|
|
parts_readings.append(f"{name}={'-' if value is None else value}")
|
|
|
|
|
|
|
|
|
|
readings = " ".join(parts_readings)
|
2026-08-10 21:37:46 -07:00
|
|
|
|
|
|
|
|
parts = []
|
|
|
|
|
for state in self.states:
|
|
|
|
|
if state.error:
|
|
|
|
|
parts.append(f"{state.rule.name}=ERROR")
|
|
|
|
|
continue
|
|
|
|
|
if not state.last_value:
|
|
|
|
|
continue
|
|
|
|
|
dwell = state.rule.dwell or 0
|
|
|
|
|
held = state.held_for(now)
|
|
|
|
|
if state.ripe(now):
|
|
|
|
|
parts.append(f"{state.rule.name}=READY")
|
|
|
|
|
else:
|
|
|
|
|
parts.append(
|
|
|
|
|
f"{state.rule.name}={format_duration(held)}/{format_duration(dwell)}"
|
|
|
|
|
)
|
|
|
|
|
|
|
|
|
|
floor = self.profile.battery_floor
|
|
|
|
|
if self.floor_latched:
|
|
|
|
|
status = f"FLOOR LATCHED until {floor.field} >= {floor.release:g}"
|
|
|
|
|
elif floor is not None and self.floor_since is not None:
|
|
|
|
|
held = now - self.floor_since
|
|
|
|
|
status = (
|
|
|
|
|
f"FLOOR ARMING {format_duration(held)}/"
|
|
|
|
|
f"{format_duration(floor.dwell)}"
|
|
|
|
|
)
|
|
|
|
|
elif parts:
|
|
|
|
|
status = "; ".join(parts)
|
|
|
|
|
else:
|
|
|
|
|
status = "no rule matches"
|
|
|
|
|
|
|
|
|
|
target = "?" if self.last_commanded is None else self._word(self.last_commanded)
|
|
|
|
|
return f"{readings} | target={target} | {status}"
|
|
|
|
|
|
|
|
|
|
def heartbeat(self, variables, now, force=False):
|
|
|
|
|
due = force or self.dry_run or (now - self.last_heartbeat) >= self.heartbeat_every
|
|
|
|
|
if not due:
|
|
|
|
|
return
|
|
|
|
|
self.last_heartbeat = now
|
|
|
|
|
self.report(self.summarize(variables, now))
|
|
|
|
|
|
2026-08-11 05:55:45 -07:00
|
|
|
async def undo_stale_action(self, session, now):
|
|
|
|
|
if self.stale_action_state is None:
|
|
|
|
|
return
|
|
|
|
|
|
|
|
|
|
previous = self.pre_stale_state
|
|
|
|
|
applied = self.stale_action_state
|
|
|
|
|
self.stale_action_state = None
|
|
|
|
|
self.pre_stale_state = None
|
|
|
|
|
|
|
|
|
|
if previous is None or previous == applied:
|
|
|
|
|
return
|
|
|
|
|
|
|
|
|
|
if self.last_commanded != applied:
|
|
|
|
|
return
|
|
|
|
|
|
|
|
|
|
self.report(
|
|
|
|
|
f"undoing the precautionary {self._word(applied)}, restoring "
|
|
|
|
|
f"{self._word(previous)} so the rules decide from here",
|
|
|
|
|
force=True,
|
|
|
|
|
)
|
|
|
|
|
await self.apply(
|
|
|
|
|
session,
|
|
|
|
|
previous,
|
|
|
|
|
"restored after telemetry recovered",
|
|
|
|
|
now,
|
|
|
|
|
force=True,
|
2026-08-12 10:53:34 -07:00
|
|
|
cause="recovered",
|
2026-08-11 05:55:45 -07:00
|
|
|
)
|
|
|
|
|
|
2026-08-10 21:37:46 -07:00
|
|
|
async def check_floor(self, session, variables, now):
|
|
|
|
|
floor = self.profile.battery_floor
|
|
|
|
|
if floor is None or not floor.enabled:
|
|
|
|
|
return False
|
|
|
|
|
|
|
|
|
|
value = variables.get(floor.field)
|
|
|
|
|
if not isinstance(value, (int, float)) or isinstance(value, bool):
|
|
|
|
|
if self.floor_latched:
|
|
|
|
|
self.report(
|
|
|
|
|
f"battery floor: {floor.field} is unreadable, holding the latch",
|
|
|
|
|
force=True,
|
|
|
|
|
)
|
|
|
|
|
return True
|
|
|
|
|
return False
|
|
|
|
|
|
|
|
|
|
if self.floor_latched:
|
|
|
|
|
if value >= floor.release:
|
|
|
|
|
self.floor_latched = False
|
|
|
|
|
self.floor_since = None
|
|
|
|
|
self.report(
|
|
|
|
|
f"battery floor released, {floor.field} back to {value:g}",
|
|
|
|
|
force=True,
|
|
|
|
|
)
|
|
|
|
|
if floor.notify_release:
|
|
|
|
|
await self.notify_floor(variables, value, released=True)
|
|
|
|
|
return False
|
|
|
|
|
|
|
|
|
|
await self.apply(
|
|
|
|
|
session,
|
|
|
|
|
floor.desired_state(),
|
|
|
|
|
f"battery floor holding, {floor.field}={value:g}",
|
|
|
|
|
now,
|
|
|
|
|
variables=variables,
|
|
|
|
|
force=True,
|
2026-08-12 10:53:34 -07:00
|
|
|
cause="floor",
|
2026-08-10 21:37:46 -07:00
|
|
|
)
|
|
|
|
|
return True
|
|
|
|
|
|
|
|
|
|
if value > floor.threshold:
|
|
|
|
|
self.floor_since = None
|
|
|
|
|
return False
|
|
|
|
|
|
|
|
|
|
if self.floor_since is None:
|
|
|
|
|
self.floor_since = now
|
|
|
|
|
|
|
|
|
|
if (now - self.floor_since) < (floor.dwell or 0):
|
|
|
|
|
return True
|
|
|
|
|
|
|
|
|
|
self.floor_latched = True
|
|
|
|
|
self.report(
|
|
|
|
|
f"BATTERY FLOOR TRIPPED: {floor.field}={value:g} at or below "
|
|
|
|
|
f"{floor.threshold:g}",
|
|
|
|
|
force=True,
|
|
|
|
|
)
|
|
|
|
|
|
|
|
|
|
await self.apply(
|
|
|
|
|
session,
|
|
|
|
|
floor.desired_state(),
|
|
|
|
|
f"battery floor, {floor.field}={value:g}",
|
|
|
|
|
now,
|
|
|
|
|
variables=variables,
|
|
|
|
|
force=True,
|
2026-08-12 10:53:34 -07:00
|
|
|
cause="floor",
|
2026-08-10 21:37:46 -07:00
|
|
|
)
|
|
|
|
|
|
|
|
|
|
if floor.notify:
|
|
|
|
|
await self.notify_floor(variables, value)
|
|
|
|
|
|
|
|
|
|
return True
|
|
|
|
|
|
|
|
|
|
async def notify_floor(self, variables, value, released=False):
|
|
|
|
|
floor = self.profile.battery_floor
|
|
|
|
|
if floor is None:
|
|
|
|
|
return
|
|
|
|
|
|
|
|
|
|
if not self.notifier.available():
|
|
|
|
|
if not released:
|
|
|
|
|
self.report(
|
|
|
|
|
"battery floor tripped but no notification channel is enabled",
|
|
|
|
|
force=True,
|
|
|
|
|
)
|
|
|
|
|
return
|
|
|
|
|
|
|
|
|
|
settings = self.profile.notifications
|
|
|
|
|
desired = None if released else floor.desired_state()
|
|
|
|
|
context = self.context(desired, None, variables, "safety")
|
|
|
|
|
context["value"] = f"{value:g}"
|
|
|
|
|
context["field"] = floor.field
|
|
|
|
|
context["threshold"] = f"{floor.threshold:g}"
|
|
|
|
|
context["release"] = f"{floor.release:g}"
|
|
|
|
|
context["reason"] = f"{floor.field} at {value:g}"
|
|
|
|
|
|
|
|
|
|
template = floor.release_template if released else floor.notify_template
|
|
|
|
|
body = render(template, context)
|
|
|
|
|
title = render(settings.title or "{profile}", context)
|
|
|
|
|
|
|
|
|
|
await self.notifier.send(
|
|
|
|
|
title,
|
|
|
|
|
body,
|
|
|
|
|
priority="urgent" if not released else None,
|
|
|
|
|
key="battery_floor_release" if released else "battery_floor",
|
|
|
|
|
force=True,
|
|
|
|
|
)
|
|
|
|
|
|
2026-08-12 10:53:34 -07:00
|
|
|
def record_event(self, desired, reason, rule, variables, cause="rule"):
|
|
|
|
|
source_identity = self.anker_profile.get("identity", {})
|
|
|
|
|
target_identity = self.target.profile.get("identity", {})
|
|
|
|
|
|
|
|
|
|
detail = {}
|
|
|
|
|
for key in (
|
|
|
|
|
"battery_soc",
|
|
|
|
|
"pv_total",
|
|
|
|
|
"pv_surplus",
|
|
|
|
|
"output_power_total",
|
|
|
|
|
"ac_input_power",
|
|
|
|
|
):
|
|
|
|
|
value = (variables or {}).get(key)
|
|
|
|
|
if isinstance(value, (int, float)) and not isinstance(value, bool):
|
|
|
|
|
detail[key] = round(value)
|
|
|
|
|
|
|
|
|
|
event = {
|
|
|
|
|
"epoch": time.time(),
|
|
|
|
|
"when": stamp(),
|
|
|
|
|
"state": bool(desired),
|
|
|
|
|
"cause": cause,
|
|
|
|
|
"profile": self.profile.name,
|
|
|
|
|
"rule": rule.name if rule else None,
|
|
|
|
|
"condition": rule.when_source if rule else None,
|
|
|
|
|
"reason": reason,
|
|
|
|
|
"target": target_identity.get("name")
|
|
|
|
|
or target_identity.get("model")
|
|
|
|
|
or self.target.host,
|
|
|
|
|
"target_channel": self.target.channel,
|
|
|
|
|
"source": source_identity.get("name")
|
|
|
|
|
or source_identity.get("model")
|
|
|
|
|
or source_identity.get("serial"),
|
|
|
|
|
"source_serial": source_identity.get("serial"),
|
|
|
|
|
"values": detail,
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
path = paths.STATE_DIR / f"events-{slugify(self.profile.name)}.json"
|
|
|
|
|
try:
|
|
|
|
|
paths.STATE_DIR.mkdir(parents=True, exist_ok=True)
|
|
|
|
|
history = []
|
|
|
|
|
if path.exists():
|
|
|
|
|
try:
|
|
|
|
|
history = json.loads(path.read_text(encoding="utf-8"))
|
|
|
|
|
except Exception:
|
|
|
|
|
history = []
|
|
|
|
|
if not isinstance(history, list):
|
|
|
|
|
history = []
|
|
|
|
|
history.append(event)
|
|
|
|
|
path.write_text(json.dumps(history[-EVENT_HISTORY:]), encoding="utf-8")
|
|
|
|
|
except Exception:
|
|
|
|
|
pass
|
|
|
|
|
|
2026-08-11 06:56:00 -07:00
|
|
|
def publish_live(self, variables, age):
|
|
|
|
|
identity = self.anker_profile.get("identity", {})
|
|
|
|
|
serial = identity.get("serial") or "unknown"
|
|
|
|
|
|
2026-08-12 11:59:04 -07:00
|
|
|
thresholds = []
|
|
|
|
|
floor = self.profile.battery_floor
|
|
|
|
|
if floor is not None and floor.enabled and floor.field == "battery_soc":
|
|
|
|
|
thresholds.append(
|
|
|
|
|
{
|
|
|
|
|
"at": floor.threshold,
|
|
|
|
|
"label": "safety floor",
|
|
|
|
|
"kind": "floor",
|
|
|
|
|
"condition": f"battery_soc <= {floor.threshold:g}",
|
|
|
|
|
"explain": (
|
|
|
|
|
f"Emergency. Turns the switch "
|
|
|
|
|
f"{'on' if floor.desired_state() else 'off'} when the "
|
|
|
|
|
f"battery reaches {floor.threshold:g}% and holds it there "
|
|
|
|
|
f"until {floor.release:g}%. Ignores the rate limits and "
|
|
|
|
|
"overrides every rule."
|
|
|
|
|
),
|
|
|
|
|
"compound": False,
|
|
|
|
|
}
|
|
|
|
|
)
|
|
|
|
|
thresholds.append(
|
|
|
|
|
{
|
|
|
|
|
"at": floor.release,
|
|
|
|
|
"label": "safety floor releases",
|
|
|
|
|
"kind": "release",
|
|
|
|
|
"condition": f"battery_soc >= {floor.release:g}",
|
|
|
|
|
"explain": (
|
|
|
|
|
f"The safety floor stops overriding the rules once the "
|
|
|
|
|
f"battery reaches {floor.release:g}%."
|
|
|
|
|
),
|
|
|
|
|
"compound": False,
|
|
|
|
|
}
|
|
|
|
|
)
|
|
|
|
|
|
|
|
|
|
for rule in self.profile.active_rules():
|
|
|
|
|
for point in battery_thresholds(rule):
|
|
|
|
|
thresholds.append(point)
|
|
|
|
|
|
2026-08-11 06:56:00 -07:00
|
|
|
payload = {
|
|
|
|
|
"serial": serial,
|
|
|
|
|
"name": identity.get("name") or identity.get("model") or serial,
|
|
|
|
|
"model": identity.get("model") or identity.get("part_number") or "",
|
2026-08-12 11:59:04 -07:00
|
|
|
"thresholds": thresholds,
|
2026-08-11 06:56:00 -07:00
|
|
|
"updated": stamp(),
|
|
|
|
|
"epoch": time.time(),
|
|
|
|
|
"age_seconds": round(age) if age is not None else None,
|
|
|
|
|
"profile": self.profile.name,
|
|
|
|
|
"target_state": self.last_commanded,
|
|
|
|
|
"floor_latched": self.floor_latched,
|
|
|
|
|
"values": {
|
|
|
|
|
key: value
|
|
|
|
|
for key, value in variables.items()
|
|
|
|
|
if isinstance(value, (int, float, bool)) or value is None
|
|
|
|
|
},
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
try:
|
|
|
|
|
paths.STATE_DIR.mkdir(parents=True, exist_ok=True)
|
|
|
|
|
destination = paths.STATE_DIR / f"live-{serial}.json"
|
|
|
|
|
destination.write_text(json.dumps(payload), encoding="utf-8")
|
|
|
|
|
except Exception:
|
|
|
|
|
pass
|
|
|
|
|
|
2026-08-12 15:38:08 -07:00
|
|
|
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
|
|
|
|
|
|
2026-08-10 21:37:46 -07:00
|
|
|
def save_state(self, desired, reason):
|
|
|
|
|
record = {
|
|
|
|
|
"profile": self.profile.name,
|
|
|
|
|
"updated": stamp(),
|
|
|
|
|
"target_state": bool(desired),
|
|
|
|
|
"reason": reason,
|
|
|
|
|
"actions_last_hour": len(self.recent_actions),
|
|
|
|
|
}
|
|
|
|
|
try:
|
|
|
|
|
paths.STATE_DIR.mkdir(parents=True, exist_ok=True)
|
|
|
|
|
existing = {}
|
|
|
|
|
if paths.RUNTIME_STATE.exists():
|
|
|
|
|
existing = json.loads(paths.RUNTIME_STATE.read_text(encoding="utf-8"))
|
|
|
|
|
existing[self.profile.name] = record
|
|
|
|
|
paths.RUNTIME_STATE.write_text(
|
|
|
|
|
json.dumps(existing, indent=2), encoding="utf-8"
|
|
|
|
|
)
|
|
|
|
|
except Exception:
|
|
|
|
|
pass
|
|
|
|
|
|
|
|
|
|
async def tick(self, session):
|
|
|
|
|
now = time.monotonic()
|
|
|
|
|
status = self.source.read()
|
|
|
|
|
age = self.source.age_seconds()
|
|
|
|
|
|
|
|
|
|
if not status:
|
|
|
|
|
self.report("no telemetry yet")
|
|
|
|
|
return
|
|
|
|
|
|
|
|
|
|
if not self.source.connected():
|
2026-08-11 05:55:45 -07:00
|
|
|
if self.disconnected_since is None:
|
|
|
|
|
self.disconnected_since = now
|
|
|
|
|
self.report(
|
|
|
|
|
"MQTT session disconnected, waiting "
|
|
|
|
|
f"{format_duration(self.profile.stale_after)} to see if it "
|
|
|
|
|
"recovers before acting",
|
|
|
|
|
force=True,
|
|
|
|
|
)
|
|
|
|
|
|
|
|
|
|
waited = now - self.disconnected_since
|
|
|
|
|
if waited < (self.profile.stale_after or 0):
|
|
|
|
|
return
|
|
|
|
|
|
2026-08-10 21:37:46 -07:00
|
|
|
if not self.stale_reported:
|
2026-08-11 05:55:45 -07:00
|
|
|
self.report(
|
|
|
|
|
f"MQTT session still down after {format_duration(waited)}, "
|
|
|
|
|
f"applying on_stale={self.profile.on_stale}",
|
|
|
|
|
force=True,
|
|
|
|
|
)
|
2026-08-10 21:37:46 -07:00
|
|
|
self.stale_reported = True
|
2026-08-11 05:55:45 -07:00
|
|
|
|
2026-08-10 21:37:46 -07:00
|
|
|
if self.profile.on_stale == "stop":
|
|
|
|
|
raise RuntimeError("stopping: MQTT session disconnected")
|
|
|
|
|
if self.profile.on_stale == "safe_state":
|
2026-08-11 05:55:45 -07:00
|
|
|
if self.stale_action_state is None:
|
|
|
|
|
self.pre_stale_state = self.last_commanded
|
|
|
|
|
self.stale_action_state = self.profile.safe_state
|
2026-08-10 21:37:46 -07:00
|
|
|
await self.apply(
|
2026-08-12 10:53:34 -07:00
|
|
|
session,
|
|
|
|
|
self.profile.safe_state,
|
|
|
|
|
"connection to Anker lost",
|
|
|
|
|
now,
|
|
|
|
|
cause="stale",
|
2026-08-10 21:37:46 -07:00
|
|
|
)
|
2026-08-11 05:59:02 -07:00
|
|
|
|
|
|
|
|
if not self.stale_notified:
|
|
|
|
|
self.stale_notified = True
|
|
|
|
|
await self.notify(
|
|
|
|
|
None,
|
|
|
|
|
None,
|
|
|
|
|
{},
|
|
|
|
|
event="stale",
|
|
|
|
|
extra={
|
|
|
|
|
"reason": "connection to Anker lost",
|
|
|
|
|
"action": self._word(self.profile.safe_state)
|
|
|
|
|
if self.profile.on_stale == "safe_state"
|
|
|
|
|
else "left as it was",
|
|
|
|
|
},
|
|
|
|
|
)
|
2026-08-10 21:37:46 -07:00
|
|
|
return
|
|
|
|
|
|
2026-08-11 05:55:45 -07:00
|
|
|
if self.disconnected_since is not None:
|
|
|
|
|
self.report(
|
|
|
|
|
"MQTT session recovered after "
|
|
|
|
|
f"{format_duration(now - self.disconnected_since)}",
|
|
|
|
|
force=True,
|
|
|
|
|
)
|
2026-08-11 05:59:02 -07:00
|
|
|
outage = format_duration(now - self.disconnected_since)
|
2026-08-11 05:55:45 -07:00
|
|
|
self.disconnected_since = None
|
|
|
|
|
await self.undo_stale_action(session, now)
|
2026-08-11 05:59:02 -07:00
|
|
|
if self.stale_notified:
|
|
|
|
|
self.stale_notified = False
|
|
|
|
|
self.stale_reported = False
|
|
|
|
|
await self.notify(
|
|
|
|
|
None, None, {}, event="recovered", extra={"outage": outage}
|
|
|
|
|
)
|
2026-08-11 05:55:45 -07:00
|
|
|
|
2026-08-10 21:37:46 -07:00
|
|
|
if age is not None and self.profile.stale_after and age > self.profile.stale_after:
|
|
|
|
|
if not self.stale_reported:
|
|
|
|
|
self.report(
|
|
|
|
|
f"telemetry stale ({format_duration(age)} old), "
|
|
|
|
|
f"policy={self.profile.on_stale}",
|
|
|
|
|
force=True,
|
|
|
|
|
)
|
|
|
|
|
self.stale_reported = True
|
|
|
|
|
|
2026-08-11 05:59:02 -07:00
|
|
|
if not self.stale_notified:
|
|
|
|
|
self.stale_notified = True
|
|
|
|
|
await self.notify(
|
|
|
|
|
None,
|
|
|
|
|
None,
|
|
|
|
|
{},
|
|
|
|
|
event="stale",
|
|
|
|
|
extra={
|
|
|
|
|
"reason": f"no message for {format_duration(age)}",
|
|
|
|
|
"action": self._word(self.profile.safe_state)
|
|
|
|
|
if self.profile.on_stale == "safe_state"
|
|
|
|
|
else "left as it was",
|
|
|
|
|
},
|
|
|
|
|
)
|
2026-08-10 21:37:46 -07:00
|
|
|
|
|
|
|
|
if self.profile.on_stale == "stop":
|
|
|
|
|
raise RuntimeError("stopping: telemetry went stale")
|
|
|
|
|
if self.profile.on_stale == "safe_state":
|
2026-08-11 05:55:45 -07:00
|
|
|
if self.stale_action_state is None:
|
|
|
|
|
self.pre_stale_state = self.last_commanded
|
|
|
|
|
self.stale_action_state = self.profile.safe_state
|
2026-08-10 21:37:46 -07:00
|
|
|
await self.apply(
|
2026-08-12 10:53:34 -07:00
|
|
|
session,
|
|
|
|
|
self.profile.safe_state,
|
|
|
|
|
"telemetry went stale",
|
|
|
|
|
now,
|
|
|
|
|
cause="stale",
|
2026-08-10 21:37:46 -07:00
|
|
|
)
|
|
|
|
|
return
|
|
|
|
|
|
|
|
|
|
if self.stale_reported:
|
|
|
|
|
self.report("telemetry recovered", force=True)
|
|
|
|
|
self.stale_reported = False
|
2026-08-11 05:55:45 -07:00
|
|
|
await self.undo_stale_action(session, now)
|
2026-08-10 21:37:46 -07:00
|
|
|
|
|
|
|
|
variables = derived_values(self.anker_profile, status)
|
2026-08-11 05:59:02 -07:00
|
|
|
self.last_variables = dict(variables)
|
|
|
|
|
self.last_seen_at = stamp()
|
2026-08-11 06:56:00 -07:00
|
|
|
self.publish_live(variables, age)
|
2026-08-12 15:38:08 -07:00
|
|
|
serial = self.anker_profile.get("identity", {}).get("serial") or "unknown"
|
|
|
|
|
self.record_telemetry(serial, variables)
|
2026-08-10 21:37:46 -07:00
|
|
|
|
|
|
|
|
missing = sorted(
|
|
|
|
|
name
|
|
|
|
|
for name in self.required_fields()
|
|
|
|
|
if name not in variables or variables[name] is None
|
|
|
|
|
)
|
|
|
|
|
if missing:
|
|
|
|
|
if not self.incomplete_reported:
|
|
|
|
|
self.report(
|
|
|
|
|
f"waiting for telemetry field(s) {missing}. Not evaluating any "
|
|
|
|
|
"rule or the safety floor until they arrive.",
|
|
|
|
|
force=True,
|
|
|
|
|
)
|
|
|
|
|
self.incomplete_reported = True
|
|
|
|
|
return
|
|
|
|
|
|
|
|
|
|
if self.incomplete_reported:
|
|
|
|
|
self.report("telemetry complete, resuming evaluation", force=True)
|
|
|
|
|
self.incomplete_reported = False
|
|
|
|
|
|
|
|
|
|
if await self.check_floor(session, variables, now):
|
|
|
|
|
self.heartbeat(variables, now)
|
|
|
|
|
return
|
|
|
|
|
|
|
|
|
|
winner, ripe = self.evaluate(variables, now)
|
|
|
|
|
|
|
|
|
|
for state in self.states:
|
|
|
|
|
if state.error:
|
|
|
|
|
self.report(f"rule {state.rule.name!r} error: {state.error}")
|
|
|
|
|
|
|
|
|
|
self.heartbeat(variables, now)
|
|
|
|
|
|
|
|
|
|
if winner is None:
|
|
|
|
|
return
|
|
|
|
|
|
|
|
|
|
desired = winner.rule.desired_state()
|
|
|
|
|
if desired is None:
|
|
|
|
|
return
|
|
|
|
|
|
|
|
|
|
reason = f"{winner.rule.name} [{winner.rule.when_source}]"
|
|
|
|
|
if len(ripe) > 1:
|
|
|
|
|
reason += f" (priority over {len(ripe) - 1} other)"
|
|
|
|
|
|
|
|
|
|
await self.apply(session, desired, reason, now, rule=winner.rule, variables=variables)
|
|
|
|
|
|
|
|
|
|
async def run(self, cycles=None):
|
|
|
|
|
await self.source.start(required=self.required_fields())
|
|
|
|
|
self.report(f"source: {self.source.label} {self.source.serial}", force=True)
|
|
|
|
|
if self.profile.battery_floor:
|
|
|
|
|
self.report(
|
|
|
|
|
f"battery floor: {self.profile.battery_floor.describe()}", force=True
|
|
|
|
|
)
|
|
|
|
|
else:
|
|
|
|
|
self.report("battery floor: NONE SET", force=True)
|
2026-08-10 22:24:01 -07:00
|
|
|
|
|
|
|
|
if self.profile.monitor_fields:
|
|
|
|
|
self.report(
|
|
|
|
|
"also logging (not used by any rule): "
|
|
|
|
|
+ ", ".join(self.profile.monitor_fields),
|
|
|
|
|
force=True,
|
|
|
|
|
)
|
2026-08-10 21:37:46 -07:00
|
|
|
self.report(f"target: {self.target.label} channel {self.target.channel}", force=True)
|
|
|
|
|
|
|
|
|
|
async with aiohttp.ClientSession() as session:
|
|
|
|
|
current = await self.target.get_state(session)
|
|
|
|
|
if current is None:
|
|
|
|
|
self.report("warning: could not read the Shelly current state", force=True)
|
|
|
|
|
else:
|
|
|
|
|
self.last_commanded = bool(current)
|
|
|
|
|
self.report(f"target currently {self._word(current)}", force=True)
|
|
|
|
|
|
|
|
|
|
try:
|
|
|
|
|
conflicts = await self.target.conflicts(session)
|
|
|
|
|
except Exception:
|
|
|
|
|
conflicts = []
|
|
|
|
|
|
|
|
|
|
if conflicts:
|
|
|
|
|
self.report("=" * 60, force=True)
|
|
|
|
|
self.report(
|
|
|
|
|
f"{len(conflicts)} CONFLICTING AUTOMATION(S) ON THE SHELLY ITSELF",
|
|
|
|
|
force=True,
|
|
|
|
|
)
|
|
|
|
|
for item in conflicts:
|
|
|
|
|
self.report(f" {item}", force=True)
|
|
|
|
|
self.report(
|
|
|
|
|
"These run on the device and will fight these rules. "
|
|
|
|
|
"Remove them in the Shelly app before relying on this.",
|
|
|
|
|
force=True,
|
|
|
|
|
)
|
|
|
|
|
self.report("=" * 60, force=True)
|
|
|
|
|
|
|
|
|
|
if self.dry_run:
|
|
|
|
|
self.report(
|
|
|
|
|
"DRY RUN: evaluating normally, but no switch command will be "
|
|
|
|
|
"sent and no notification will fire",
|
|
|
|
|
force=True,
|
|
|
|
|
)
|
|
|
|
|
|
|
|
|
|
count = 0
|
|
|
|
|
try:
|
|
|
|
|
while cycles is None or count < cycles:
|
|
|
|
|
await self.tick(session)
|
|
|
|
|
count += 1
|
|
|
|
|
if cycles is not None and count >= cycles:
|
|
|
|
|
self.report(
|
|
|
|
|
f"completed {count} cycle(s), stopping as requested",
|
|
|
|
|
force=True,
|
|
|
|
|
)
|
|
|
|
|
break
|
|
|
|
|
await interruptible_sleep(self.profile.poll_interval)
|
|
|
|
|
finally:
|
|
|
|
|
await self.source.stop()
|
|
|
|
|
await asyncio.sleep(0.25)
|
|
|
|
|
|
|
|
|
|
async def close(self):
|
|
|
|
|
await self.source.stop()
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
async def dry_run_report(profile, cycles=3, offline=False, overrides=None):
|
|
|
|
|
problems, notes = validate(profile)
|
|
|
|
|
|
|
|
|
|
print()
|
|
|
|
|
print(f"Power profile: {profile.name}")
|
|
|
|
|
print(f" file {paths.relative(profile.path)}")
|
|
|
|
|
print(f" source {profile.source_reference}")
|
|
|
|
|
print(f" target {profile.target_reference}")
|
|
|
|
|
print(f" enabled {profile.enabled}")
|
|
|
|
|
print()
|
|
|
|
|
|
|
|
|
|
for note in notes:
|
|
|
|
|
print(f" note: {note}")
|
|
|
|
|
for problem in problems:
|
|
|
|
|
print(f" PROBLEM: {problem}")
|
|
|
|
|
|
|
|
|
|
if problems:
|
|
|
|
|
print()
|
|
|
|
|
print(f"{len(problems)} problem(s) found. Fix these before running.")
|
|
|
|
|
return False
|
|
|
|
|
|
|
|
|
|
print(f" syntax OK, {len(profile.active_rules())} active rule(s)")
|
|
|
|
|
|
|
|
|
|
settings = profile.notifications
|
|
|
|
|
if settings.enabled:
|
|
|
|
|
from .notify import Notifier
|
|
|
|
|
|
|
|
|
|
probe = Notifier(settings)
|
|
|
|
|
channels = probe.available()
|
|
|
|
|
missing = probe.missing()
|
|
|
|
|
print(
|
|
|
|
|
f" notifications on via {', '.join(channels) if channels else 'NO CHANNEL'}"
|
|
|
|
|
f", throttle {format_duration(settings.throttle)}"
|
|
|
|
|
)
|
|
|
|
|
if missing:
|
|
|
|
|
print(f" requested but not enabled: {', '.join(missing)}")
|
|
|
|
|
if not channels:
|
|
|
|
|
print(" nothing will be delivered until a channel is enabled")
|
|
|
|
|
else:
|
|
|
|
|
print(" notifications off")
|
|
|
|
|
|
|
|
|
|
if profile.battery_floor:
|
|
|
|
|
floor = profile.battery_floor
|
|
|
|
|
print(f" safety floor: {floor.describe()}, dwell {format_duration(floor.dwell)}")
|
|
|
|
|
else:
|
|
|
|
|
print(" safety floor: NONE SET")
|
|
|
|
|
|
|
|
|
|
if overrides:
|
|
|
|
|
print()
|
|
|
|
|
print("Simulated overrides:")
|
|
|
|
|
for key, value in sorted(overrides.items()):
|
|
|
|
|
print(f" {key} = {value!r}")
|
|
|
|
|
|
|
|
|
|
if offline:
|
|
|
|
|
anker_profile = load_yaml(profile.source_path)
|
|
|
|
|
samples = {
|
|
|
|
|
key: spec.get("sample")
|
|
|
|
|
for key, spec in (anker_profile.get("readable") or {}).items()
|
|
|
|
|
}
|
|
|
|
|
variables = derived_values(anker_profile, samples)
|
|
|
|
|
if overrides:
|
|
|
|
|
variables.update(overrides)
|
|
|
|
|
print()
|
|
|
|
|
print("Offline evaluation against the sample values in the device profile:")
|
|
|
|
|
referenced = sorted(profile.referenced_names())
|
|
|
|
|
overridden = overrides or {}
|
|
|
|
|
if referenced:
|
|
|
|
|
readings = ", ".join(
|
|
|
|
|
f"{name}={variables.get(name)!r}"
|
|
|
|
|
+ (" (simulated)" if name in overridden else "")
|
|
|
|
|
for name in referenced
|
|
|
|
|
)
|
|
|
|
|
print(f" values: {readings}")
|
|
|
|
|
|
|
|
|
|
floor_wins = _print_floor_verdict(profile, variables)
|
|
|
|
|
|
|
|
|
|
print()
|
|
|
|
|
if floor_wins:
|
|
|
|
|
print(" Rules below are shown for reference, but the floor would")
|
|
|
|
|
print(" take precedence while it is latched:")
|
|
|
|
|
_print_rule_table(profile, variables, overrides=overrides, skip_values=True)
|
|
|
|
|
print()
|
|
|
|
|
print("Sample values are a snapshot from discovery, not live data.")
|
|
|
|
|
return True
|
|
|
|
|
|
|
|
|
|
print()
|
|
|
|
|
print(f"Connecting for a live dry run ({cycles} cycle(s), no commands sent)...")
|
|
|
|
|
|
|
|
|
|
engine = Engine(profile, dry_run=True)
|
|
|
|
|
await engine.source.start(required=engine.required_fields())
|
|
|
|
|
|
|
|
|
|
try:
|
|
|
|
|
async with aiohttp.ClientSession() as session:
|
|
|
|
|
reachable = await engine.target.reachable(session)
|
|
|
|
|
current = await engine.target.get_state(session)
|
|
|
|
|
print()
|
|
|
|
|
print(f" Shelly reachable: {reachable}")
|
|
|
|
|
if current is not None:
|
|
|
|
|
print(f" Shelly currently: {'ON' if current else 'OFF'}")
|
|
|
|
|
if not reachable:
|
|
|
|
|
print(
|
|
|
|
|
" the target did not respond; check access.host in its profile"
|
|
|
|
|
)
|
|
|
|
|
|
|
|
|
|
for index in range(cycles):
|
|
|
|
|
if index:
|
|
|
|
|
await interruptible_sleep(profile.poll_interval)
|
|
|
|
|
|
|
|
|
|
status = engine.source.read()
|
|
|
|
|
if not status:
|
|
|
|
|
print()
|
|
|
|
|
print(f" cycle {index + 1}: no telemetry decoded yet")
|
|
|
|
|
continue
|
|
|
|
|
|
|
|
|
|
variables = derived_values(engine.anker_profile, status)
|
|
|
|
|
if overrides:
|
|
|
|
|
variables.update(overrides)
|
|
|
|
|
|
|
|
|
|
absent = sorted(
|
|
|
|
|
name
|
|
|
|
|
for name in engine.required_fields()
|
|
|
|
|
if name not in variables or variables[name] is None
|
|
|
|
|
)
|
|
|
|
|
if absent:
|
|
|
|
|
print()
|
|
|
|
|
print(f" cycle {index + 1}: waiting for field(s) {absent}")
|
|
|
|
|
print(" nothing is evaluated until they arrive")
|
|
|
|
|
continue
|
|
|
|
|
|
|
|
|
|
now = time.monotonic()
|
|
|
|
|
for state in engine.states:
|
|
|
|
|
state.update(variables, now)
|
|
|
|
|
|
|
|
|
|
print()
|
|
|
|
|
print(f" cycle {index + 1} (telemetry age "
|
|
|
|
|
f"{format_duration(engine.source.age_seconds())})")
|
|
|
|
|
floor_wins = _print_floor_verdict(profile, variables)
|
|
|
|
|
print()
|
|
|
|
|
if floor_wins:
|
|
|
|
|
print(" Rules below are shown for reference, but the floor")
|
|
|
|
|
print(" would take precedence while it is latched:")
|
|
|
|
|
_print_rule_table(
|
|
|
|
|
profile, variables, engine, now, overrides=overrides
|
|
|
|
|
)
|
|
|
|
|
finally:
|
|
|
|
|
await engine.source.stop()
|
|
|
|
|
|
|
|
|
|
print()
|
|
|
|
|
print("Dry run complete. No commands were sent.")
|
|
|
|
|
return True
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def _preview_notification(profile, rule, variables, desired):
|
|
|
|
|
from .notify import render
|
|
|
|
|
|
|
|
|
|
settings = profile.notifications
|
|
|
|
|
if not settings.wants("action") or not settings.rule_enabled(rule):
|
|
|
|
|
return None
|
|
|
|
|
|
|
|
|
|
source_identity = load_yaml(profile.source_path).get("identity", {})
|
|
|
|
|
target_profile = load_yaml(profile.target_path)
|
|
|
|
|
target_identity = target_profile.get("identity", {})
|
|
|
|
|
|
|
|
|
|
context = dict(variables)
|
|
|
|
|
context.update(
|
|
|
|
|
{
|
|
|
|
|
"profile": profile.name,
|
|
|
|
|
"rule": rule.name,
|
|
|
|
|
"condition": rule.when_source,
|
|
|
|
|
"action": "ON" if desired else "OFF",
|
|
|
|
|
"action_word": "on" if desired else "off",
|
|
|
|
|
"source_name": (
|
|
|
|
|
source_identity.get("name")
|
|
|
|
|
or source_identity.get("model")
|
|
|
|
|
or source_identity.get("serial")
|
|
|
|
|
),
|
|
|
|
|
"source_model": source_identity.get("model", ""),
|
|
|
|
|
"source_serial": source_identity.get("serial", ""),
|
|
|
|
|
"target_name": (
|
|
|
|
|
target_identity.get("name")
|
|
|
|
|
or target_identity.get("model")
|
|
|
|
|
or target_profile.get("access", {}).get("host")
|
|
|
|
|
),
|
|
|
|
|
"target_model": target_identity.get("model", ""),
|
|
|
|
|
"target_host": target_profile.get("access", {}).get("host", ""),
|
|
|
|
|
"target_channel": profile.target_channel or 0,
|
|
|
|
|
"time": stamp(),
|
|
|
|
|
"event": "action",
|
|
|
|
|
}
|
|
|
|
|
)
|
|
|
|
|
return render(settings.template_for(rule), context)
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def _print_floor_verdict(profile, variables):
|
|
|
|
|
floor = profile.battery_floor
|
|
|
|
|
if floor is None or not floor.enabled:
|
|
|
|
|
return False
|
|
|
|
|
|
|
|
|
|
value = variables.get(floor.field)
|
|
|
|
|
|
|
|
|
|
print()
|
|
|
|
|
if not isinstance(value, (int, float)) or isinstance(value, bool):
|
|
|
|
|
print(f" SAFETY FLOOR: {floor.field} is not readable, cannot evaluate")
|
|
|
|
|
return False
|
|
|
|
|
|
|
|
|
|
if value <= floor.threshold:
|
|
|
|
|
print(
|
|
|
|
|
f" SAFETY FLOOR TRIPS: {floor.field}={value:g} is at or below "
|
|
|
|
|
f"{floor.threshold:g}"
|
|
|
|
|
)
|
|
|
|
|
print(
|
|
|
|
|
f" after {format_duration(floor.dwell)} it would {floor.action} "
|
|
|
|
|
f"and LATCH until {floor.field} reaches {floor.release:g}"
|
|
|
|
|
)
|
|
|
|
|
print(" it bypasses the rate limits and outranks every rule below")
|
|
|
|
|
print(" while latched, no rule can turn the target off")
|
|
|
|
|
return True
|
|
|
|
|
|
|
|
|
|
print(
|
|
|
|
|
f" safety floor idle: {floor.field}={value:g} is above "
|
|
|
|
|
f"{floor.threshold:g}"
|
|
|
|
|
)
|
|
|
|
|
return False
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def _print_rule_table(
|
|
|
|
|
profile, variables, engine=None, now=None, overrides=None, skip_values=False
|
|
|
|
|
):
|
|
|
|
|
referenced = sorted(profile.referenced_names())
|
|
|
|
|
overrides = overrides or {}
|
|
|
|
|
if referenced and not skip_values:
|
|
|
|
|
readings = ", ".join(
|
|
|
|
|
f"{name}={variables.get(name)!r}"
|
|
|
|
|
+ (" (simulated)" if name in overrides else "")
|
|
|
|
|
for name in referenced
|
|
|
|
|
)
|
|
|
|
|
print(f" values: {readings}")
|
|
|
|
|
|
|
|
|
|
states = engine.states if engine else None
|
|
|
|
|
|
|
|
|
|
for index, rule in enumerate(profile.active_rules()):
|
2026-08-12 09:53:46 -07:00
|
|
|
acting = False
|
|
|
|
|
errored = False
|
|
|
|
|
|
2026-08-10 21:37:46 -07:00
|
|
|
if states:
|
|
|
|
|
state = states[index]
|
|
|
|
|
satisfied = state.last_value
|
|
|
|
|
held = state.held_for(now) if now else 0
|
|
|
|
|
ripe = state.ripe(now) if now else False
|
|
|
|
|
if state.error:
|
|
|
|
|
verdict = f"ERROR {state.error}"
|
2026-08-12 09:53:46 -07:00
|
|
|
errored = True
|
2026-08-10 21:37:46 -07:00
|
|
|
elif not satisfied:
|
2026-08-12 09:53:46 -07:00
|
|
|
verdict = "no"
|
2026-08-10 21:37:46 -07:00
|
|
|
elif ripe:
|
2026-08-12 09:53:46 -07:00
|
|
|
verdict = f"YES, ready now -> would {rule.action}"
|
|
|
|
|
acting = True
|
2026-08-10 21:37:46 -07:00
|
|
|
else:
|
|
|
|
|
remaining = (rule.dwell or 0) - held
|
2026-08-12 09:53:46 -07:00
|
|
|
verdict = (
|
|
|
|
|
f"YES, holding {format_duration(remaining)} more before acting"
|
|
|
|
|
)
|
|
|
|
|
acting = True
|
2026-08-10 21:37:46 -07:00
|
|
|
else:
|
|
|
|
|
try:
|
|
|
|
|
satisfied = rule.evaluate(variables)
|
2026-08-12 09:53:46 -07:00
|
|
|
if satisfied:
|
|
|
|
|
verdict = (
|
|
|
|
|
f"YES -> would {rule.action} after "
|
|
|
|
|
f"{format_duration(rule.dwell)}"
|
|
|
|
|
)
|
|
|
|
|
acting = True
|
|
|
|
|
else:
|
|
|
|
|
verdict = "no"
|
2026-08-10 21:37:46 -07:00
|
|
|
except Exception as err:
|
|
|
|
|
verdict = f"ERROR {type(err).__name__}: {err}"
|
2026-08-12 09:53:46 -07:00
|
|
|
errored = True
|
2026-08-10 21:37:46 -07:00
|
|
|
|
2026-08-12 09:53:46 -07:00
|
|
|
marker = " *" if acting else (" " if not errored else " !")
|
|
|
|
|
print(f" {marker} [{rule.priority:>3}] {rule.name}")
|
2026-08-10 21:37:46 -07:00
|
|
|
print(f" when {rule.when_source}")
|
|
|
|
|
print(f" {verdict}")
|
|
|
|
|
|
|
|
|
|
desired = rule.desired_state()
|
2026-08-12 09:53:46 -07:00
|
|
|
if acting and desired is not None:
|
2026-08-10 21:37:46 -07:00
|
|
|
if engine:
|
|
|
|
|
message = _render_live_notification(engine, rule, variables, desired)
|
|
|
|
|
else:
|
|
|
|
|
message = _preview_notification(profile, rule, variables, desired)
|
|
|
|
|
if message:
|
2026-08-12 09:53:46 -07:00
|
|
|
print(f" would notify: {message}")
|
2026-08-10 21:37:46 -07:00
|
|
|
|
|
|
|
|
|
|
|
|
|
def _render_live_notification(engine, rule, variables, desired):
|
|
|
|
|
from .notify import render
|
|
|
|
|
|
|
|
|
|
settings = engine.profile.notifications
|
|
|
|
|
if not settings.wants("action") or not settings.rule_enabled(rule):
|
|
|
|
|
return None
|
|
|
|
|
context = engine.context(desired, rule, variables, "action")
|
|
|
|
|
return render(settings.template_for(rule), context)
|