737 lines
24 KiB
Python
737 lines
24 KiB
Python
import asyncio
|
|
import time
|
|
from pathlib import Path
|
|
|
|
from aiohttp import ClientSession
|
|
|
|
from anker_solix_api.api import AnkerSolixApi
|
|
from anker_solix_api.mqtt_factory import SolixMqttDeviceFactory
|
|
|
|
from . import paths
|
|
from .credentials import load_credentials
|
|
from .profiles import classify, load_yaml, now_iso, save_yaml, slugify
|
|
|
|
MODEL_NAMES = {
|
|
"A1782": "SOLIX F3000",
|
|
"A1790": "SOLIX F3800",
|
|
"A1780": "SOLIX F2000",
|
|
"A1781": "SOLIX F2600",
|
|
"A1761": "SOLIX C1000X",
|
|
"A1753": "SOLIX C800",
|
|
"A17C1": "Solarbank 2",
|
|
"A17C5": "Solarbank 3",
|
|
"A17C0": "Solarbank E1600",
|
|
"A17X8": "Smart Plug",
|
|
"A2345": "Prime Charger",
|
|
"AS200": "Alternator Charger",
|
|
}
|
|
|
|
REALTIME_TIMEOUT = 60
|
|
POLLER_TIMEOUT = 60
|
|
IGNORED_KEYS = {"topics"}
|
|
|
|
ONLINE_KEYS = ("wifi_online", "is_online", "online", "wifi_connected")
|
|
|
|
BLUETOOTH_ONLY_HINT = """ This device reports as not cloud connected.
|
|
|
|
Some models (for example the F2000 / PowerHouse 767) pair with the Anker
|
|
app over Bluetooth only. You press a button on the unit to make it
|
|
discoverable, and it has no persistent cloud connection. Those devices
|
|
publish nothing to Anker's MQTT broker, which refuses the subscription.
|
|
|
|
Such a device cannot be used as an automation source here. Local
|
|
Bluetooth access needs a different project entirely (SolixBLE).
|
|
|
|
If you believe this device IS cloud connected, force an attempt with
|
|
--include-offline"""
|
|
|
|
OFFLINE_HINT = """ Check that the device is:
|
|
1. powered on and awake, not in standby
|
|
2. connected to WiFi, not only paired over Bluetooth
|
|
3. showing as online in the Anker app right now
|
|
This tool reads from Anker's cloud MQTT broker, so the device must be
|
|
reachable by the cloud. A Bluetooth-only connection is not enough."""
|
|
|
|
PROFILE_HEADER = """
|
|
Anker SOLIX device profile.
|
|
|
|
Generated by: solixauto discover-anker
|
|
This file is READ-ONLY input for the automation engine. The engine never
|
|
sends control commands to Anker devices; it only reads the fields below.
|
|
|
|
readable: fields observed in live MQTT telemetry, with the type and a
|
|
sample value captured at generation time.
|
|
derived: computed fields available to power-profile rules.
|
|
writable: control methods the library exposes for this model. Listed for
|
|
reference only. The automation engine does not call them.
|
|
|
|
Field names ending in ? or starting with unknown_ are not yet confirmed by
|
|
the upstream project. Avoid building rules on them.
|
|
|
|
Regenerate this file after a library update to pick up new fields.
|
|
"""
|
|
|
|
|
|
def model_label(part_number):
|
|
friendly = MODEL_NAMES.get(str(part_number).upper())
|
|
return f"{part_number} ({friendly})" if friendly else str(part_number)
|
|
|
|
|
|
def part_number_of(info):
|
|
return str((info or {}).get("device_pn") or (info or {}).get("product_code") or "")
|
|
|
|
|
|
def device_label(info):
|
|
return (
|
|
(info or {}).get("device_name")
|
|
or (info or {}).get("name")
|
|
or (info or {}).get("alias_name")
|
|
or (info or {}).get("alias")
|
|
or ""
|
|
)
|
|
|
|
|
|
def online_hint(info):
|
|
for key in ONLINE_KEYS:
|
|
if key in (info or {}):
|
|
value = info[key]
|
|
if isinstance(value, str):
|
|
lowered = value.strip().lower()
|
|
if lowered in ("false", "0", "no", "offline"):
|
|
return False
|
|
if lowered in ("true", "1", "yes", "online"):
|
|
return True
|
|
continue
|
|
if value is None:
|
|
continue
|
|
return bool(value)
|
|
return None
|
|
|
|
|
|
BLUETOOTH_ONLY_SHORT = (
|
|
"Bluetooth-only models cannot be used as a source. "
|
|
"Details: --explain-skipped"
|
|
)
|
|
|
|
|
|
def clean_status(raw):
|
|
return {key: value for key, value in (raw or {}).items() if key not in IGNORED_KEYS}
|
|
|
|
|
|
def build_derived(status):
|
|
derived = {}
|
|
|
|
pv_keys = sorted(k for k in status if k.startswith("pv_") and k.endswith("_power"))
|
|
if pv_keys:
|
|
derived["pv_total"] = {
|
|
"expression": " + ".join(pv_keys),
|
|
"description": "collective solar input in watts",
|
|
}
|
|
|
|
usb_keys = sorted(
|
|
k
|
|
for k in status
|
|
if (k.startswith("usbc_") or k.startswith("usba_")) and k.endswith("_power")
|
|
)
|
|
if usb_keys:
|
|
derived["usb_total"] = {
|
|
"expression": " + ".join(usb_keys),
|
|
"description": "combined USB output in watts",
|
|
}
|
|
|
|
if pv_keys and "output_power_total" in status:
|
|
derived["pv_surplus"] = {
|
|
"expression": " + ".join(pv_keys) + " - output_power_total",
|
|
"description": (
|
|
"solar minus load in watts. Positive means the sun is covering "
|
|
"everything drawing from the unit and the surplus charges the "
|
|
"battery. Negative means the battery is making up the shortfall."
|
|
),
|
|
}
|
|
derived["pv_covers_load"] = {
|
|
"expression": " + ".join(pv_keys) + " >= output_power_total",
|
|
"description": "true when solar alone is carrying the load",
|
|
}
|
|
|
|
if "ac_input_power" in status and "output_power_total" in status:
|
|
derived["net_power"] = {
|
|
"expression": "ac_input_power - output_power_total",
|
|
"description": "positive when importing, negative when discharging",
|
|
}
|
|
|
|
if "battery_soc" in status and "min_soc" in status:
|
|
derived["soc_headroom"] = {
|
|
"expression": "battery_soc - min_soc",
|
|
"description": "percentage points above the configured floor",
|
|
}
|
|
|
|
return derived
|
|
|
|
|
|
def introspect_writable(device):
|
|
if device is None:
|
|
return {}
|
|
|
|
writable = {}
|
|
for name in sorted(dir(device)):
|
|
if not name.startswith("set_"):
|
|
continue
|
|
attribute = getattr(device, name, None)
|
|
if not callable(attribute):
|
|
continue
|
|
entry = {"method": name}
|
|
try:
|
|
import inspect
|
|
|
|
signature = inspect.signature(attribute)
|
|
params = [
|
|
p for p in signature.parameters if p not in ("self", "args", "kwargs")
|
|
]
|
|
if params:
|
|
entry["parameters"] = params
|
|
except (TypeError, ValueError):
|
|
pass
|
|
writable[name[4:]] = entry
|
|
|
|
commands = getattr(device, "commands", None)
|
|
if isinstance(commands, dict):
|
|
for name in sorted(commands):
|
|
writable.setdefault(str(name), {"command": str(name)})
|
|
|
|
return writable
|
|
|
|
|
|
def find_duplicates(serial, keep):
|
|
import yaml
|
|
|
|
duplicates = []
|
|
if not paths.ANKER_PROFILE_DIR.exists() or not serial:
|
|
return duplicates
|
|
|
|
wanted = str(serial).lower()
|
|
|
|
for candidate in sorted(paths.ANKER_PROFILE_DIR.iterdir()):
|
|
if candidate.suffix not in (".yaml", ".yml") or candidate == keep:
|
|
continue
|
|
try:
|
|
with candidate.open("r", encoding="utf-8") as handle:
|
|
data = yaml.safe_load(handle)
|
|
except Exception:
|
|
continue
|
|
if not isinstance(data, dict):
|
|
continue
|
|
|
|
identity = data.get("identity") or {}
|
|
if str(identity.get("serial") or "").lower() == wanted:
|
|
duplicates.append(candidate)
|
|
|
|
return duplicates
|
|
|
|
|
|
def profiles_referencing(path):
|
|
import yaml
|
|
|
|
referencing = []
|
|
if not paths.POWER_PROFILE_DIR.exists():
|
|
return referencing
|
|
|
|
names = {path.name.lower(), path.stem.lower()}
|
|
|
|
for candidate in sorted(paths.POWER_PROFILE_DIR.iterdir()):
|
|
if candidate.suffix not in (".yaml", ".yml"):
|
|
continue
|
|
try:
|
|
with candidate.open("r", encoding="utf-8") as handle:
|
|
data = yaml.safe_load(handle)
|
|
except Exception:
|
|
continue
|
|
if not isinstance(data, dict):
|
|
continue
|
|
|
|
source = (data.get("source") or {}).get("profile")
|
|
if source and str(source).lower() in names:
|
|
referencing.append(candidate)
|
|
|
|
return referencing
|
|
|
|
|
|
def build_profile(device_sn, info, status, device):
|
|
part_number = part_number_of(info) or "unknown"
|
|
|
|
readable = {}
|
|
for key in sorted(status):
|
|
value = status[key]
|
|
readable[key] = {"type": classify(value), "sample": value}
|
|
|
|
name = device_label(info)
|
|
aliases = [entry for entry in (name, device_sn, part_number) if entry]
|
|
|
|
return {
|
|
"kind": "anker",
|
|
"generated": now_iso(),
|
|
"aliases": aliases,
|
|
"identity": {
|
|
"serial": device_sn,
|
|
"part_number": part_number,
|
|
"model": MODEL_NAMES.get(part_number.upper(), "unknown"),
|
|
"name": name,
|
|
"make": "Anker SOLIX",
|
|
},
|
|
"access": {
|
|
"transport": "anker-cloud-mqtt",
|
|
"engine_mode": "read-only",
|
|
"realtime_trigger_timeout_seconds": REALTIME_TIMEOUT,
|
|
},
|
|
"derived": build_derived(status),
|
|
"readable": readable,
|
|
"writable": introspect_writable(device),
|
|
}
|
|
|
|
|
|
class MqttReader:
|
|
def __init__(self, api, session, device_sn, info):
|
|
self.api = api
|
|
self.session = session
|
|
self.device_sn = device_sn
|
|
self.info = info
|
|
self.device = None
|
|
self.topics = set()
|
|
self.poller = None
|
|
self.last_message = 0.0
|
|
self.message_count = 0
|
|
|
|
def _callback(self, session, topic, message, data, model, *args, **kwargs):
|
|
self.last_message = time.monotonic()
|
|
self.message_count += 1
|
|
|
|
def build_topics(self):
|
|
topics = set()
|
|
prefix = self.session.get_topic_prefix(deviceDict=self.info)
|
|
if prefix:
|
|
topics.add(f"{prefix}#")
|
|
command_prefix = self.session.get_topic_prefix(
|
|
deviceDict=self.info, publish=True
|
|
)
|
|
if command_prefix:
|
|
topics.add(f"{command_prefix}#")
|
|
return topics
|
|
|
|
def prepare(self):
|
|
self.topics = self.build_topics()
|
|
if not self.topics:
|
|
raise RuntimeError(
|
|
f"could not resolve an MQTT topic prefix for {self.device_sn}. "
|
|
"The device may not be owned by this account."
|
|
)
|
|
self.device = SolixMqttDeviceFactory(self.api, self.device_sn).create_device()
|
|
return self.topics
|
|
|
|
async def start(self, realtime=True):
|
|
self.prepare()
|
|
|
|
trigger_devices = {self.device_sn} if realtime else set()
|
|
|
|
self.poller = asyncio.get_running_loop().create_task(
|
|
self.session.message_poller(
|
|
topics=self.topics,
|
|
trigger_devices=trigger_devices,
|
|
msg_callback=self._callback,
|
|
timeout=POLLER_TIMEOUT,
|
|
)
|
|
)
|
|
|
|
await asyncio.sleep(2)
|
|
self.request_status()
|
|
|
|
def request_status(self):
|
|
try:
|
|
result = self.session.status_request(deviceDict=self.info, wait_for_publish=2)
|
|
return bool(result and result.is_published())
|
|
except Exception:
|
|
return False
|
|
|
|
def read(self):
|
|
data = getattr(self.session, "mqtt_data", None) or {}
|
|
return clean_status(data.get(self.device_sn) or {})
|
|
|
|
async def wait_for_data(self, seconds=45, verbose=False, required=None):
|
|
deadline = time.monotonic() + seconds
|
|
requested_again = False
|
|
required = set(required or ())
|
|
|
|
while time.monotonic() < deadline:
|
|
status = self.read()
|
|
if status and not (required - set(status)):
|
|
return status
|
|
|
|
if self.poller and self.poller.done():
|
|
error = self.poller.exception()
|
|
if error:
|
|
raise RuntimeError(f"MQTT poller stopped: {error}")
|
|
|
|
remaining = deadline - time.monotonic()
|
|
if not requested_again and remaining < seconds / 2:
|
|
requested_again = True
|
|
self.request_status()
|
|
if verbose:
|
|
print(" no data yet, re-requesting status...")
|
|
|
|
await asyncio.sleep(1)
|
|
|
|
return self.read()
|
|
|
|
async def settle(self, seconds=8):
|
|
best = self.read()
|
|
deadline = time.monotonic() + seconds
|
|
while time.monotonic() < deadline:
|
|
await asyncio.sleep(2)
|
|
latest = self.read()
|
|
if len(latest) > len(best):
|
|
best = latest
|
|
return best
|
|
|
|
def age_seconds(self):
|
|
if not self.last_message:
|
|
return None
|
|
return time.monotonic() - self.last_message
|
|
|
|
async def stop(self):
|
|
if self.poller is not None:
|
|
self.poller.cancel()
|
|
try:
|
|
await self.poller
|
|
except asyncio.CancelledError:
|
|
pass
|
|
except Exception:
|
|
pass
|
|
self.poller = None
|
|
|
|
|
|
async def open_session(verbose=True):
|
|
user, password, country = load_credentials()
|
|
websession = ClientSession()
|
|
api = AnkerSolixApi(user, password, country, websession, None)
|
|
|
|
try:
|
|
if await api.async_authenticate():
|
|
if verbose:
|
|
print("Anker cloud authentication: OK")
|
|
elif verbose:
|
|
print("Anker cloud authentication: cached token")
|
|
|
|
await api.update_sites()
|
|
await api.get_bind_devices()
|
|
|
|
mqtt_session = await api.startMqttSession()
|
|
if not mqtt_session:
|
|
raise RuntimeError("startMqttSession returned nothing")
|
|
if not mqtt_session.is_connected():
|
|
raise RuntimeError("MQTT session did not connect")
|
|
|
|
if verbose:
|
|
print(f"Connected to MQTT server {mqtt_session.host}:{mqtt_session.port}")
|
|
|
|
return api, websession, mqtt_session
|
|
except Exception:
|
|
await websession.close()
|
|
raise
|
|
|
|
|
|
async def close_session(api, websession):
|
|
session = getattr(api, "mqttsession", None)
|
|
if session is not None:
|
|
cleanup = getattr(session, "cleanup", None)
|
|
if callable(cleanup):
|
|
try:
|
|
cleanup()
|
|
except Exception:
|
|
pass
|
|
await websession.close()
|
|
|
|
|
|
async def discover(
|
|
settle=45,
|
|
only_pn=None,
|
|
only_sn=None,
|
|
skip=None,
|
|
include_offline=False,
|
|
explain_skipped=False,
|
|
verbose=True,
|
|
):
|
|
paths.ensure_dirs()
|
|
|
|
if verbose:
|
|
print("Authenticating and enumerating devices...")
|
|
|
|
api, websession, mqtt_session = await open_session(verbose=verbose)
|
|
written = []
|
|
|
|
try:
|
|
devices = api.devices or {}
|
|
if not devices:
|
|
raise RuntimeError("no owned devices returned for this account")
|
|
|
|
skip = {s.strip() for s in (skip or []) if s.strip()}
|
|
|
|
targets = {}
|
|
for serial, info in devices.items():
|
|
if only_sn and serial != only_sn:
|
|
continue
|
|
if only_pn and part_number_of(info).upper() != only_pn.upper():
|
|
continue
|
|
if serial in skip:
|
|
if verbose:
|
|
print(f" skipping {serial} (--skip)")
|
|
continue
|
|
|
|
if not include_offline and online_hint(info) is False:
|
|
label = model_label(part_number_of(info))
|
|
print(f" {serial} - {label}: not cloud connected, skipped.")
|
|
if explain_skipped:
|
|
print(BLUETOOTH_ONLY_HINT)
|
|
else:
|
|
print(f" {BLUETOOTH_ONLY_SHORT}")
|
|
continue
|
|
|
|
targets[serial] = info
|
|
|
|
if not targets:
|
|
raise RuntimeError("no devices matched, or all matches were offline")
|
|
|
|
if verbose:
|
|
print(f"Found {len(targets)} device(s) to harvest.")
|
|
|
|
readers = {}
|
|
all_topics = set()
|
|
|
|
for serial, info in targets.items():
|
|
reader = MqttReader(api, mqtt_session, serial, info)
|
|
try:
|
|
all_topics |= reader.prepare()
|
|
readers[serial] = reader
|
|
except Exception as err:
|
|
print(f" {serial}: {err}")
|
|
|
|
if not readers:
|
|
raise RuntimeError("no device topics could be resolved")
|
|
|
|
def shared_callback(session, topic, message, data, model, *args, **kwargs):
|
|
for entry in readers.values():
|
|
if topic and entry.topics:
|
|
for pattern in entry.topics:
|
|
if topic.startswith(pattern.rstrip("#")):
|
|
entry._callback(
|
|
session, topic, message, data, model, *args, **kwargs
|
|
)
|
|
return
|
|
|
|
if verbose:
|
|
print(
|
|
f"Subscribing {len(all_topics)} topic(s) for "
|
|
f"{len(readers)} device(s) on one session..."
|
|
)
|
|
|
|
poller = asyncio.get_running_loop().create_task(
|
|
mqtt_session.message_poller(
|
|
topics=all_topics,
|
|
trigger_devices=set(readers),
|
|
msg_callback=shared_callback,
|
|
timeout=POLLER_TIMEOUT,
|
|
)
|
|
)
|
|
|
|
try:
|
|
await asyncio.sleep(3)
|
|
|
|
for reader in readers.values():
|
|
reader.request_status()
|
|
|
|
for serial, reader in readers.items():
|
|
info = targets[serial]
|
|
label = model_label(part_number_of(info))
|
|
if verbose:
|
|
print()
|
|
print(f" {serial} - {label}")
|
|
print(f" waiting up to {settle}s for telemetry...")
|
|
|
|
status = await reader.wait_for_data(settle, verbose=verbose)
|
|
|
|
if not status:
|
|
print(
|
|
f" no telemetry decoded after {settle}s "
|
|
f"({reader.message_count} message(s) seen)"
|
|
)
|
|
if reader.message_count:
|
|
print(
|
|
" messages arrived but decoded to nothing. This model "
|
|
"may not have field mappings in mqttmap.py yet."
|
|
)
|
|
else:
|
|
print(" no messages received from this device.")
|
|
print(OFFLINE_HINT)
|
|
print()
|
|
print(
|
|
" Once it is online, retry just this device with:"
|
|
)
|
|
print(f" {paths.command('discover-anker --sn ' + serial)}")
|
|
continue
|
|
|
|
status = await reader.settle(8)
|
|
|
|
profile = build_profile(serial, info, status, reader.device)
|
|
friendly = profile["identity"].get("name") or ""
|
|
if friendly:
|
|
stem = slugify(friendly).lower()
|
|
else:
|
|
stem = (
|
|
f"{slugify(part_number_of(info) or 'device')}-"
|
|
f"{slugify(serial)}"
|
|
).lower()
|
|
destination = paths.ANKER_PROFILE_DIR / f"{stem}.yaml"
|
|
|
|
stale = find_duplicates(serial, destination)
|
|
|
|
save_yaml(destination, profile, header=PROFILE_HEADER)
|
|
written.append(destination)
|
|
|
|
for other in stale:
|
|
print()
|
|
print(
|
|
f" WARNING: {other.name} also describes this device."
|
|
)
|
|
users = profiles_referencing(other)
|
|
if users:
|
|
names = ", ".join(p.name for p in users)
|
|
print(f" {names} still points at the old file.")
|
|
print(
|
|
f" Change its source to {destination.name}, "
|
|
"then remove:"
|
|
)
|
|
else:
|
|
print(" Two profiles for one device is confusing.")
|
|
print(" Remove the old one:")
|
|
print(f" rm {other}")
|
|
|
|
if verbose:
|
|
print(
|
|
f" {len(profile['readable'])} readable, "
|
|
f"{len(profile['derived'])} derived, "
|
|
f"{len(profile['writable'])} control method(s)"
|
|
)
|
|
print(f" wrote {paths.relative(destination)}")
|
|
finally:
|
|
poller.cancel()
|
|
try:
|
|
await poller
|
|
except asyncio.CancelledError:
|
|
pass
|
|
except Exception:
|
|
pass
|
|
finally:
|
|
await close_session(api, websession)
|
|
|
|
return written
|
|
|
|
|
|
class AnkerSource:
|
|
def __init__(self, profile_path):
|
|
self.profile_path = Path(profile_path)
|
|
self.profile = load_yaml(self.profile_path)
|
|
identity = self.profile.get("identity", {})
|
|
self.serial = identity.get("serial")
|
|
self.part_number = identity.get("part_number")
|
|
self.label = model_label(self.part_number)
|
|
self.derived = self.profile.get("derived", {}) or {}
|
|
|
|
if not self.serial:
|
|
raise ValueError(f"{self.profile_path} has no identity.serial")
|
|
|
|
self._api = None
|
|
self._websession = None
|
|
self._mqtt = None
|
|
self._reader = None
|
|
|
|
def raw_requirements(self, required):
|
|
from .rules import expression_names
|
|
|
|
derived = self.profile.get("derived") or {}
|
|
readable = set((self.profile.get("readable") or {}).keys())
|
|
|
|
resolved = set()
|
|
for name in required or ():
|
|
if name in derived:
|
|
spec = derived[name]
|
|
expression = spec.get("expression") if isinstance(spec, dict) else spec
|
|
try:
|
|
resolved |= expression_names(expression)
|
|
except Exception:
|
|
pass
|
|
elif not readable or name in readable:
|
|
resolved.add(name)
|
|
|
|
return {name for name in resolved if name not in derived}
|
|
|
|
async def start(self, settle=45, required=None):
|
|
required = self.raw_requirements(required)
|
|
self._api, self._websession, self._mqtt = await open_session(verbose=False)
|
|
|
|
info = (self._api.devices or {}).get(self.serial)
|
|
if info is None:
|
|
raise RuntimeError(f"serial {self.serial} is not owned by this account")
|
|
|
|
self._reader = MqttReader(self._api, self._mqtt, self.serial, info)
|
|
await self._reader.start(realtime=True)
|
|
try:
|
|
status = await self._reader.wait_for_data(settle, required=required)
|
|
except Exception:
|
|
await self.stop()
|
|
raise
|
|
|
|
if not status:
|
|
await self.stop()
|
|
raise RuntimeError(
|
|
f"no telemetry from {self.label} {self.serial} after {settle}s.\n"
|
|
+ OFFLINE_HINT
|
|
)
|
|
|
|
missing = set(required or ()) - set(status)
|
|
if missing:
|
|
await self.stop()
|
|
raise RuntimeError(
|
|
f"telemetry from {self.serial} is missing field(s) "
|
|
f"{sorted(missing)} after {settle}s. Check the field names in your "
|
|
"power profile against: solixauto fields <device>"
|
|
)
|
|
|
|
def read(self):
|
|
if self._reader is None:
|
|
return {}
|
|
return self._reader.read()
|
|
|
|
def age_seconds(self):
|
|
if self._reader is None:
|
|
return None
|
|
return self._reader.age_seconds()
|
|
|
|
def connected(self):
|
|
if self._mqtt is None:
|
|
return False
|
|
try:
|
|
return bool(self._mqtt.is_connected())
|
|
except Exception:
|
|
return False
|
|
|
|
async def trigger(self):
|
|
if self._reader is None:
|
|
return False
|
|
return self._reader.request_status()
|
|
|
|
async def stop(self):
|
|
if self._reader is not None:
|
|
await self._reader.stop()
|
|
self._reader = None
|
|
if self._api is not None and self._websession is not None:
|
|
await close_session(self._api, self._websession)
|
|
self._api = None
|
|
self._websession = None
|