import json import threading import time import urllib.error import urllib.request from collections import deque from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer from . import paths from .profiles import list_profiles, load_yaml, save_yaml, slugify HISTORY_SIZE = 2880 POLL_TIMEOUT = 3 class Sampler: def __init__(self, interval=5.0, history=HISTORY_SIZE): self.interval = interval self.history = deque(maxlen=history) self.snapshot = { "anker": [], "shelly": [], "events": [], "multi": False, "updated": None, "engine": False, } self.lock = threading.Lock() self.stop_event = threading.Event() self.thread = None def start(self): self.thread = threading.Thread(target=self._loop, daemon=True) self.thread.start() def stop(self): self.stop_event.set() def _loop(self): while not self.stop_event.is_set(): try: self._sample() except Exception: pass self.stop_event.wait(self.interval) def _read_live(self, serial): path = paths.STATE_DIR / f"live-{serial}.json" if not path.exists(): return None try: data = json.loads(path.read_text(encoding="utf-8")) except Exception: return None if time.time() - float(data.get("epoch", 0)) > 120: data["fresh"] = False else: data["fresh"] = True return data def rename(self, kind, path_name, new_name, on_device=False): directory = ( paths.ANKER_PROFILE_DIR if kind == "anker" else paths.SHELLY_PROFILE_DIR ) path = directory / path_name if path.parent.resolve() != directory.resolve() or not path.exists(): raise ValueError("unknown device") new_name = str(new_name).strip() if not new_name: raise ValueError("name cannot be empty") if len(new_name) > 60: raise ValueError("name is too long") profile = load_yaml(path) identity = profile.setdefault("identity", {}) old_name = identity.get("name") or path.stem identity["name"] = new_name aliases = list(profile.get("aliases") or []) for candidate in (new_name, old_name, path.stem): if candidate and candidate not in aliases: aliases.append(candidate) profile["aliases"] = aliases destination = directory / f"{slugify(new_name).lower()}.yaml" if destination != path and destination.exists(): raise ValueError(f"{destination.name} already exists") save_yaml(path, profile) if destination != path: path.replace(destination) wrote_device = False if on_device and kind == "shelly": wrote_device = self._write_shelly_name(profile, new_name) self._sample() return {"file": destination.name, "name": new_name, "on_device": wrote_device} def _write_shelly_name(self, profile, new_name): access = profile.get("access") or {} identity = profile.get("identity") or {} host = access.get("host") generation = int(identity.get("generation", 1) or 1) if not host or generation < 2: return False payload = json.dumps({"config": {"device": {"name": new_name}}}).encode() request = urllib.request.Request( f"http://{host}/rpc/Sys.SetConfig", data=payload, headers={"Content-Type": "application/json"}, ) try: with urllib.request.urlopen(request, timeout=POLL_TIMEOUT) as response: return response.status == 200 except Exception: return False def _anker_devices(self): devices = [] for path in list_profiles(paths.ANKER_PROFILE_DIR): try: profile = load_yaml(path) except Exception: continue identity = profile.get("identity", {}) serial = identity.get("serial") live = self._read_live(serial) if serial else None values = (live or {}).get("values", {}) devices.append( { "file": path.name, "name": identity.get("name") or identity.get("model") or path.stem, "model": identity.get("model") or identity.get("part_number") or "", "serial": serial, "live": bool(live and live.get("fresh")), "updated": (live or {}).get("updated"), "age": (live or {}).get("age_seconds"), "profile": (live or {}).get("profile"), "thresholds": (live or {}).get("thresholds") or [], "floor_latched": (live or {}).get("floor_latched", False), "battery_soc": values.get("battery_soc"), "output_watts": values.get("output_power_total"), "ac_in_watts": values.get("ac_input_power"), "pv_watts": values.get("pv_total"), "pv_surplus": values.get("pv_surplus"), "temperature": values.get("temperature"), } ) return devices def _poll_shelly(self, host, generation, channel): path = "/rpc/Shelly.GetStatus" if generation >= 2 else "/status" try: with urllib.request.urlopen( f"http://{host}{path}", timeout=POLL_TIMEOUT ) as response: data = json.loads(response.read().decode("utf-8")) except Exception: return None, None if generation >= 2: entry = data.get(f"switch:{channel}") or {} return entry.get("output"), entry.get("apower") relays = data.get("relays") or [] meters = data.get("meters") or [] state = relays[channel].get("ison") if channel < len(relays) else None power = meters[channel].get("power") if channel < len(meters) else None return state, power def _utilized_shelly_files(self): utilized = set() for path in list_profiles(paths.POWER_PROFILE_DIR): try: power_profile = load_yaml(path) except Exception: continue target = (power_profile.get("target") or {}).get("profile") if not target: continue resolved = paths.resolve_profile(target, "shelly") if resolved is not None: utilized.add(resolved.name) return utilized def _shelly_devices(self): devices = [] utilized = self._utilized_shelly_files() for path in list_profiles(paths.SHELLY_PROFILE_DIR): if path.name not in utilized: continue try: profile = load_yaml(path) except Exception: continue identity = profile.get("identity", {}) access = profile.get("access", {}) host = access.get("host") generation = int(identity.get("generation", 1) or 1) channels = profile.get("channels") or {"0": {}} channel = int(sorted(channels)[0]) state, power = self._poll_shelly(host, generation, channel) if host else (None, None) devices.append( { "file": path.name, "name": identity.get("name") or path.stem, "model": identity.get("model") or "", "host": host, "channel": channel, "state": state, "watts": power, "reachable": state is not None, } ) return devices def _events(self): events = [] if not paths.STATE_DIR.exists(): return events for path in sorted(paths.STATE_DIR.glob("events-*.json")): try: data = json.loads(path.read_text(encoding="utf-8")) except Exception: continue if isinstance(data, list): events.extend(entry for entry in data if isinstance(entry, dict)) events.sort(key=lambda entry: entry.get("epoch", 0), reverse=True) return events[:60] def _sample(self): anker = self._anker_devices() shelly = self._shelly_devices() events = self._events() point = {"t": time.time()} for device in anker: if not device["live"]: continue serial = device["serial"] point[f"{serial}:soc"] = device["battery_soc"] point[f"{serial}:out"] = device["output_watts"] point[f"{serial}:ac"] = device["ac_in_watts"] point[f"{serial}:pv"] = device["pv_watts"] with self.lock: if len(point) > 1: self.history.append(point) self.snapshot = { "anker": anker, "shelly": shelly, "events": events, "multi": len(anker) > 1 or len(shelly) > 1, "updated": time.strftime("%H:%M:%S"), "engine": any(d["live"] for d in anker), "interval": self.interval, } def payload(self): with self.lock: return { **self.snapshot, "history": list(self.history), } PAGE = """