0.2.8: Removed threading to reduce errors, made MQTT threaded, PVOutput updated to not miss uploads and error

This commit is contained in:
bohdan-s
2022-01-19 15:58:22 +11:00
parent c6476aa41d
commit b5bc64ecae
4 changed files with 71 additions and 84 deletions
+24 -41
View File
@@ -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")
+45 -42
View File
@@ -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
self.last_run = time.time()
+1
View File
@@ -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")
+1 -1
View File
@@ -1 +1 @@
__version__ = '0.2.7'
__version__ = '0.2.8'