diff --git a/SunGather/exports/mqtt.py b/SunGather/exports/mqtt.py index 3b01e65..d1decba 100644 --- a/SunGather/exports/mqtt.py +++ b/SunGather/exports/mqtt.py @@ -4,7 +4,6 @@ import paho.mqtt.client as mqtt class export_mqtt(object): def __init__(self): - self._isConfigured = False self.mqtt_client = None self.sensor_topic = None self.homeassistant = False @@ -13,10 +12,14 @@ class export_mqtt(object): self.model = None self.model_clean = None self.inverter_ip = None + self.mqtt_queue = [] # Configure MQTT def configure(self, config, config_inverter): self.mqtt_client = mqtt.Client(client_id="SunGather") + self.mqtt_client.on_connect = self.on_connect + self.mqtt_client.on_disconnect = self.on_disconnect + self.mqtt_client.on_publish = self.on_publish if config.get('username') and config.get('password'): self.mqtt_client.username_pw_set(config.get('username'), config.get('password')) @@ -24,16 +27,10 @@ class export_mqtt(object): if config.get('port') == 8883: self.mqtt_client.tls_set() - try: - self.mqtt_client.connect(config.get('host'), port=config.get('port', 1883), keepalive=60) - self._isConfigured = True - except Exception as err: - logging.error(f"MQTT: Connection {config.get('host')}:{config.get('port', 1883)}") - logging.error(f"MQTT: Error: {err}") - return False + self.mqtt_client.connect_async(config.get('host'), port=config.get('port', 1883), keepalive=60) + self.mqtt_client.loop_start() self.inverter_ip = config_inverter.get('host') - self.sensor_topic = config.get('topic', 'inverter/{model}/registers') self.homeassistant = config.get('homeassistant', False) @@ -41,37 +38,32 @@ class export_mqtt(object): for ha_sensor in config.get('ha_sensors'): self.ha_sensors.append(ha_sensor) - logging.info(f"MQTT: Configured {config.get('host')}:{config.get('port', 1883)}, HA Enabled = {config.get('homeassistant', False)}") + return True - return self._isConfigured + def on_connect(self, client, userdata, flags, rc): + logging.info(f"MQTT: Connected {client._host}:{client._port}") + + def on_disconnect(self, client, userdata, rc): + logging.info(f"MQTT: Server Disconnected code:{rc}") + + def on_publish(self, client, userdata, mid): + self.mqtt_queue.remove(mid) + logging.info(f"MQTT: Message {mid} Published") def publish(self, inverter): - global mqtt_client - - if not self._isConfigured: - logging.info("MQTT: Skipped, Initial Configuration Failed") - return False + if not self.mqtt_client.is_connected(): + logging.warning(f'MQTT: Server Disconnected; {self.mqtt_queue.__len__()} messages queued') + elif self.mqtt_queue.__len__() > 10: + logging.warning(f'MQTT: {self.mqtt_queue.__len__()} messages queued, this may be due to a MQTT server issue') if not self.model: self.model = inverter.get('device_type_code', 'unknown') self.model_clean = self.model.replace('.','').replace('-','') self.sensor_topic = self.sensor_topic.replace('{model}', self.model_clean) - logging.debug(f'MQTT: Sensor Topic = {self.sensor_topic}') - + logging.debug(f'MQTT: Sensor Topic = {self.sensor_topic}') logging.debug(f"MQTT: Publishing: {self.sensor_topic} : {json.dumps(inverter)}") - try: - result = self.mqtt_client.publish(self.sensor_topic, json.dumps(inverter).replace('"', '\"')) - result.wait_for_publish() - if result.rc == mqtt.MQTT_ERR_SUCCESS: - logging.info("MQTT: Published") - else: - logging.error(f"MQTT: Failed to publish with error code: {result.rc}") - if result.rc == mqtt.MQTT_ERR_NO_CONN: - logging.warning(f"MQTT: Attempting to reconnect to MQTT") - self.mqtt_client.reconnect() - except Exception as err: - logging.error(f"MQTT: Failed to publish with error: {err}") + self.mqtt_queue.append(self.mqtt_client.publish(self.sensor_topic, json.dumps(inverter).replace('"', '\"'), qos=1).mid) if self.homeassistant: if self.ha_discovery: @@ -96,17 +88,8 @@ class export_mqtt(object): config_msg['ic'] = "mdi:solar-power" config_msg['device'] = { "name":"Solar Inverter", "mf":"Sungrow", "mdl":self.model, "connections":[["address", self.inverter_ip ]]} - try: - logging.debug(f'MQTT: Topic; {ha_topic}, Message: {config_msg}') - result = self.mqtt_client.publish(ha_topic, json.dumps(config_msg), retain=True) - result.wait_for_publish() - if result.rc != mqtt.MQTT_ERR_SUCCESS: - logging.error(f"MQTT: Failed to publish with error code: {result.rc}") - if result.rc == mqtt.MQTT_ERR_NO_CONN: - logging.warning(f"MQTT: Attempting to reconnect to MQTT") - self.mqtt_client.reconnect() - except Exception as err: - logging.error(f"MQTT: Failed to publish with error: {err}") + logging.debug(f'MQTT: Topic; {ha_topic}, Message: {config_msg}') + self.mqtt_queue.append(self.mqtt_client.publish(ha_topic, json.dumps(config_msg), retain=True, qos=1).mid) self.ha_discovery = False logging.info("MQTT: Published Home Assistant Discovery messages") diff --git a/SunGather/exports/pvoutput.py b/SunGather/exports/pvoutput.py index 9523c5f..a225e98 100644 --- a/SunGather/exports/pvoutput.py +++ b/SunGather/exports/pvoutput.py @@ -1,7 +1,7 @@ import logging import requests import datetime -import json +import time # Configure PVOutput class export_pvoutput(object): @@ -16,7 +16,7 @@ class export_pvoutput(object): self.url_getsystem = self.url_base + "getsystem.jsp" self.rate_limit = None self.parameters = [] - self.latest_run = None + self.last_run = None self.tid = '1618' @property @@ -37,18 +37,18 @@ class export_pvoutput(object): self.collected_data = {} self.payload_data = None self.batch_count = 0 + self.last_run = 0 for parameter in config.get('parameters'): self.parameters.append(parameter) try: logging.debug(f"PVOutput: Get System ; {self.url_getsystem}, {str(self.headers)}, 'teams': '1'") - response = requests.post(url=self.url_getsystem,headers=self.headers, params={'teams': '1'}) + response = requests.post(url=self.url_getsystem,headers=self.headers, params={'teams': '1'}, timeout=3) logging.debug(f"PVOutput: Response; {str(response.status_code)} Message; {str(response.content)}") - if response.status_code == 200: - retload = response.content.decode() - system = retload.split(';')[0] - teams = retload.split(';')[2] + if response.status_code == 200: + system = response.text.split(';')[0] + teams = response.text.split(';')[2] invertername = system.split(',')[0] self.status_interval = int(system.split(',')[15]) @@ -61,28 +61,25 @@ class export_pvoutput(object): logging.error(f"PVOutput: System Status Failed; {str(response.status_code)} Message; {str(response.content)}") except Exception as err: - logging.error(f"PVOutput: Failed to configure with error: {err}") + logging.error(f"PVOutput: Failed to configure") + logging.debug(f"{err}") return False - join_team = config.get('join_team',True) try: + join_team = config.get('join_team',True) if not team_member and join_team: logging.debug(f"PVOutput: Join Team; {self.url_jointeam}, {str(self.headers)}, 'tid': '{self.tid}'") - response = requests.post(url=self.url_jointeam,headers=self.headers, params={'tid': self.tid}) + response = requests.post(url=self.url_jointeam,headers=self.headers, params={'tid': self.tid}, timeout=3) logging.debug(f"PVOutput: Response; {str(response.status_code)} Message; {str(response.content)}") elif team_member and not join_team: logging.debug(f"PVOutput: Leave Team; {self.url_leaveteam}, {str(self.headers)}, 'tid': '{self.tid}'") - response = requests.post(url=self.url_leaveteam,headers=self.headers, params={'tid': self.tid}) + response = requests.post(url=self.url_leaveteam,headers=self.headers, params={'tid': self.tid}, timeout=3) logging.debug(f"PVOutput: Response; {str(response.status_code)} Message; {str(response.content)}") except Exception as err: - logging.error(f"PVOutput: Team Membership update failed: {err}") - - if not response.status_code == 200: - logging.error(f"PVOutput: Team update Failed; {str(response.status_code)} Message; {str(response.content)}") + logging.debug(f"{err}") logging.info(f"PVOutput: Configured export to {invertername} every {self.status_interval} minutes") - self._isConfigured = True - return self._isConfigured + return True def publish(self, inverter): """ @@ -112,26 +109,31 @@ class export_pvoutput(object): m1 Text Message 1 No text 30 chars max Yes """ + # Check all required registers have been returned by the inverter + if not inverter.get('timestamp', False): + logging.warning('PVOutput: Skipping, Timestamp missing from last scrape') + return False + for parameter in self.parameters: + if parameter.get('name')[0] == 'v': + if not inverter.get(parameter.get('register', False)): + logging.warning(f"PVOutput: Skipping, {parameter.get('name')} configured to use {parameter.get('register')} but inverter is not returning this register") + return False + now = datetime.datetime.strptime(inverter.get('timestamp'), "%Y-%m-%d %H:%M:%S") - if not self.latest_run: # Set last run to 1 min ago if never run, that way we don't miss any uploads - self.latest_run = now - datetime.timedelta(minutes=1) # Add new data to old data and increase count of data points for parameter in self.parameters: - if parameter.get('name')[0] == 'v': + if parameter.get('name')[0] == 'v': if self.collected_data.get(parameter.get('name'),False): - if parameter.get('multiple'): - self.collected_data[parameter.get('name')] = round(self.collected_data[parameter.get('name')] + (inverter.get(parameter.get('register')) * parameter.get('multiple')),3) - else: - self.collected_data[parameter.get('name')] = round(self.collected_data[parameter.get('name')] + inverter.get(parameter.get('register')),3) + if parameter.get('multiple'): + self.collected_data[parameter.get('name')] = round(self.collected_data[parameter.get('name')] + (inverter.get(parameter.get('register')) * parameter.get('multiple')),3) + else: + self.collected_data[parameter.get('name')] = round(self.collected_data[parameter.get('name')] + inverter.get(parameter.get('register')),3) else: - if inverter.get(parameter.get('register')): if parameter.get('multiple'): self.collected_data[parameter.get('name')] = round(inverter.get(parameter.get('register')) * parameter.get('multiple'),3) else: self.collected_data[parameter.get('name')] = inverter.get(parameter.get('register')) - else: - logging.warning(f"PVOutput: {parameter.get('name')} configured to use {parameter.get('register')} but inverter is not returning this register") elif parameter.get('name') == 'c1': cumulative_energy = parameter.get('value') @@ -140,11 +142,10 @@ class export_pvoutput(object): else: self.collected_data['count'] = 1 - logging.info('PVOutput: Data logged') - logging.debug(f'PVOutput: Data: {self.collected_data}') + logging.debug(f'PVOutput: Data Logged: {self.collected_data}') # Process data points every status_interval - if int(now.strftime("%M")) % self.status_interval == 0 and not now.strftime("%H:%M") == self.latest_run.strftime("%H:%M"): + if((time.time() - self.last_run) > (self.status_interval * 60)): for v in self.collected_data: if v[0] == 'v' and not self.collected_data[v] == 0: if v[1] == '6' or v[1] == '7': # Round to 1 decimal place @@ -172,25 +173,27 @@ class export_pvoutput(object): self.batch_count +=1 - if self.batch_count == self.batch_points: - + if self.batch_count >= self.batch_points: payload = {} payload['data'] = self.payload_data if 'cumulative_energy' in locals(): payload['c1'] = cumulative_energy - logging.debug("PVOutput: Request; " + self.url_addbatchstatus + ", " + str(self.headers) + " : " + str(payload)) + try: + logging.debug("PVOutput: Request; " + self.url_addbatchstatus + ", " + str(self.headers) + " : " + str(payload)) + response = requests.post(url=self.url_addbatchstatus, headers=self.headers, params=payload, timeout=3) + self.batch_count = 0 - response = requests.post(url=self.url_addbatchstatus, headers=self.headers, params=payload) - self.batch_count = 0 - - if response.status_code != requests.codes.ok: - raise RuntimeError(response.text) - else: - self.payload_data = None - logging.info("PVOutput: Data uploaded") + if response.status_code != requests.codes.ok: + logging.error(f"PVOutput: Upload Failed; {str(response.status_code)} Message; {str(response.content)}") + else: + self.payload_data = None + logging.info("PVOutput: Data uploaded") + except Exception as err: + logging.error(f"PVOutput: Failed to Upload") + logging.debug(f"{err}") else: logging.info("PVOutput: Data added to next batch upload") - self.latest_run = now \ No newline at end of file + self.last_run = time.time() \ No newline at end of file diff --git a/SunGather/exports/webserver.py b/SunGather/exports/webserver.py index 6dd449f..da3e7ef 100644 --- a/SunGather/exports/webserver.py +++ b/SunGather/exports/webserver.py @@ -16,6 +16,7 @@ class export_webserver(object): try: self.webServer = HTTPServer(('', config.get('port',8080)), MyServer) self.t = Thread(target=self.webServer.serve_forever) + self.t.daemon = True # Make it a deamon, so if main loop ends the webserver dies self.t.start() self.config_inverter = config_inverter logging.info(f"Webserver: Configured") diff --git a/SunGather/version.py b/SunGather/version.py index 94abaec..93de59e 100644 --- a/SunGather/version.py +++ b/SunGather/version.py @@ -1 +1 @@ -__version__ = '0.2.7' \ No newline at end of file +__version__ = '0.2.8' \ No newline at end of file