Remove use of external mqtt file, just use paho library

This commit is contained in:
DSTO\pivatom
2019-01-11 11:40:04 +10:30
parent b61da848f4
commit acc287d572

View File

@@ -1,26 +1,31 @@
import Messaging.mqttsession as ms
import paho.mqtt.client as mqtt
import time
import umsgpack
class Commander:
currentVote = None
# Stores voters that connect to maintain a majority.
# Voters who do not vote in latest round are removed.
connectedVoters = {}
votes = []
taking_votes = False
def __init__(self, timeout = 60):
'''
Initial/default waiting time is 1 minute for votes to come in.
'''
self.timeout = timeout
self._votes = []
ms.Client()
self.timeout = timeout
def on_message(client, userdata, message):
self._votes.append(umsgpack.unpackb(message.payload))
ms.client.subscribe("FakeSwarm/FirstTest", on_message)
self.client = mqtt.Client()
self.client.connect('172.16.13.128')
ms.client.loop_start()
# Commander needs a will message too, for the decentralised version, so the
# voters know to pick a new commander.
def make_decision(self):
# Could change this to be a strategy, for different implementations of making
# a decision on the votes.
votes = self._votes
dif_votes = {}
@@ -40,9 +45,11 @@ class Commander:
return max_vote
def get_votes(self, topic, message, qos=0):
def get_votes():
message = {"type": "voteRequest"}
message_packed = umsgpack.packb(message)
# Publish a message that votes are needed.
ms.client.publish(topic, message, qos)
ms.client.publish("swarm1/voters", message_packed, qos)
time.sleep(self.timeout)
self.make_decision()
@@ -50,4 +57,22 @@ class Commander:
ms.client.subscribe(topic)
if callback is not None:
ms.client.message_callback_add(topic, callback)
ms.client.message_callback_add(topic, callback)
def on_message(self, client, userdata, message):
messageDict = umsgpack.unpackb()
if "type" in messageDict.keys:
if messageDict["type"] == "connect":
# Voter just connected/reconnnected.
elif messageDict["type"] == "vote":
# Voter is sending in their vote.
if messageDict[""]
else:
# Not a message we should be using.
NotImplementedError
def on_connect(client, userdata, flags, rc):
self.client.subscribe("swarm1/commander")