forked from platypush/platypush
Use the Redis backend to dispatch messages to the core bus if available
This commit is contained in:
parent
7ab85b4cfa
commit
8c208c0028
1 changed files with 11 additions and 1 deletions
|
@ -5,6 +5,7 @@ import paho.mqtt.publish as publisher
|
|||
|
||||
from platypush.backend import Backend
|
||||
from platypush.message import Message
|
||||
from platypush.context import get_backend
|
||||
|
||||
|
||||
class MqttBackend(Backend):
|
||||
|
@ -26,6 +27,11 @@ class MqttBackend(Backend):
|
|||
self.port = port
|
||||
self.topic = '{}/{}'.format(topic, self.device_id)
|
||||
|
||||
try:
|
||||
self.redis = get_backend('redis')
|
||||
except:
|
||||
self.redis = None
|
||||
|
||||
|
||||
def send_message(self, msg):
|
||||
publisher.single(self.topic, str(msg), hostname=self.host, port=self.port)
|
||||
|
@ -43,6 +49,10 @@ class MqttBackend(Backend):
|
|||
pass
|
||||
|
||||
self.logger.info('Received message on the MQTT backend: {}'.format(msg))
|
||||
|
||||
if self.redis:
|
||||
self.redis.send_message(msg)
|
||||
else:
|
||||
self.bus.post(msg)
|
||||
|
||||
super().run()
|
||||
|
|
Loading…
Reference in a new issue