Separate mqtt client in a different package

master
esensar 2018-05-08 17:13:57 +02:00
parent b85e90cfb4
commit a006a6ac5f
4 changed files with 57 additions and 38 deletions

View File

@ -26,9 +26,11 @@ def setup_blueprints(app):
from .devices import devices_bp from .devices import devices_bp
from .accounts import accounts_bp from .accounts import accounts_bp
from .api import api_bp from .api import api_bp
from .mqtt import mqtt_bp
app.register_blueprint(devices_bp) app.register_blueprint(devices_bp)
app.register_blueprint(accounts_bp) app.register_blueprint(accounts_bp)
app.register_blueprint(mqtt_bp)
app.register_blueprint(api_bp, url_prefix='/api') app.register_blueprint(api_bp, url_prefix='/api')

View File

@ -1,26 +1,11 @@
import atexit import sys
from flask import Blueprint from flask import Blueprint
from .mqtt_client import MqttClient
from .models import Device, Recording from .models import Device, Recording
from app import app
devices_bp = Blueprint('devices', __name__) devices_bp = Blueprint('devices', __name__)
# When app dies, stop mqtt connection
def on_stop():
MqttClient.tear_down()
atexit.register(on_stop)
# Setup
@devices_bp.record
def __on_blueprint_setup(setup_state):
print('Blueprint setup')
MqttClient.setup(setup_state.app)
# Public interface # Public interface
def create_device(name, device_type=1): def create_device(name, device_type=1):
""" """
@ -69,3 +54,33 @@ def get_device(device_id):
raise ValueError("Device with id %s does not exist" % device_id) raise ValueError("Device with id %s does not exist" % device_id)
return Device.get(id=device_id) return Device.get(id=device_id)
def create_recording(device_id, raw_json):
"""
Tries to create recording with given parameters. Raises error on failure
:param device_id: Id of device
:type device_id: int
:param raw_json: Raw json received
:type raw_json: json
:raises: ValueError if parsing fails or device does not exist
"""
def parse_raw_json_recording(device_id, json_msg) -> Recording:
try:
return Recording(device_id=device_id,
record_type=json_msg["record_type"],
record_value=json_msg["record_value"],
recorded_at=json_msg["recorded_at"],
raw_json=json_msg)
except KeyError:
error_type, error_instance, traceback = sys.exc_info()
raise ValueError("JSON parsing failed! Key error: "
+ str(error_instance))
if not Device.exists(id=device_id):
raise ValueError("Device does not exist!")
recording = parse_raw_json_recording(device_id, raw_json)
with app.app_context():
recording.save()

View File

@ -0,0 +1,19 @@
import atexit
from flask import Blueprint
from .mqtt_client import MqttClient
mqtt_bp = Blueprint('mqtt', __name__)
# Setup
@mqtt_bp.record
def __on_blueprint_setup(setup_state):
MqttClient.setup(setup_state.app)
# When app dies, stop mqtt connection
def on_stop():
MqttClient.tear_down()
atexit.register(on_stop)

View File

@ -1,8 +1,7 @@
import sys import sys
import json import json
from flask_mqtt import Mqtt from flask_mqtt import Mqtt
from .models import Recording import app.devices as devices
from app import app
class MqttClient: class MqttClient:
@ -49,10 +48,9 @@ class MqttClient:
print("Payload: " + message.payload.decode()) print("Payload: " + message.payload.decode())
try: try:
# If type is JSON # If type is JSON
recording = MqttClient.parse_json_message( devices.create_recording(
message.topic, message.payload.decode()) MqttClient.get_device_id(message.topic),
with app.app_context(): json.loads(message.payload.decode()))
recording.save()
except ValueError: except ValueError:
print("ERROR!") print("ERROR!")
error_type, error_instance, traceback = sys.exc_info() error_type, error_instance, traceback = sys.exc_info()
@ -60,21 +58,6 @@ class MqttClient:
print("Instance: " + str(error_instance)) print("Instance: " + str(error_instance))
return return
@staticmethod
def parse_json_message(topic, payload) -> Recording:
try:
json_msg = json.loads(payload)
device_id = MqttClient.get_device_id(topic)
return Recording(device_id=device_id,
record_type=json_msg["record_type"],
record_value=json_msg["record_value"],
recorded_at=json_msg["recorded_at"],
raw_json=json_msg)
except KeyError:
error_type, error_instance, traceback = sys.exc_info()
raise ValueError("JSON parsing failed! Key error: "
+ str(error_instance))
@staticmethod @staticmethod
def get_device_id(topic) -> int: def get_device_id(topic) -> int:
device_token, device_id = topic.split("/") device_token, device_id = topic.split("/")