Files

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