2022-04-15 20:02:42 +02:00
|
|
|
import traceback
|
|
|
|
|
2023-10-03 17:23:29 +02:00
|
|
|
from asgiref.sync import async_to_sync
|
|
|
|
from channels.generic.websocket import WebsocketConsumer
|
2022-04-04 01:13:48 +02:00
|
|
|
|
2022-04-04 14:48:43 +02:00
|
|
|
from c3nav.mesh import messages
|
2023-10-03 17:23:29 +02:00
|
|
|
from c3nav.mesh.models import MeshNode, NodeMessage
|
2022-04-04 01:13:48 +02:00
|
|
|
|
|
|
|
|
2023-10-03 17:23:29 +02:00
|
|
|
# noinspection PyAttributeOutsideInit
|
|
|
|
class MeshConsumer(WebsocketConsumer):
|
|
|
|
def connect(self):
|
2022-04-06 17:25:46 +02:00
|
|
|
print('connected!')
|
2022-04-15 20:57:11 +02:00
|
|
|
# todo: auth
|
2023-10-03 17:23:29 +02:00
|
|
|
self.uplink_node = None
|
|
|
|
self.dst_nodes = set()
|
|
|
|
self.accept()
|
2022-04-04 14:48:43 +02:00
|
|
|
|
2023-10-03 17:23:29 +02:00
|
|
|
def disconnect(self, close_code):
|
2022-04-06 17:25:46 +02:00
|
|
|
print('disconnected!')
|
2023-10-03 17:23:29 +02:00
|
|
|
if self.uplink_node is not None:
|
|
|
|
self.remove_route(self.uplink_node)
|
|
|
|
self.channel_layer.group_discard('route_%s' % self.node.address.replace(':', ''), self.channel_name)
|
|
|
|
self.channel_layer.group_discard('route_broadcast', self.channel_name)
|
2022-04-04 01:13:48 +02:00
|
|
|
|
2023-10-03 17:23:29 +02:00
|
|
|
def send_msg(self, msg):
|
2022-04-06 22:56:08 +02:00
|
|
|
print('Sending message:', msg)
|
2023-10-03 17:23:29 +02:00
|
|
|
self.send(bytes_data=msg.encode())
|
2022-04-06 22:56:08 +02:00
|
|
|
|
2023-10-03 17:23:29 +02:00
|
|
|
def receive(self, text_data=None, bytes_data=None):
|
2022-04-04 14:48:43 +02:00
|
|
|
if bytes_data is None:
|
|
|
|
return
|
2022-04-15 20:02:42 +02:00
|
|
|
try:
|
|
|
|
msg = messages.Message.decode(bytes_data)
|
|
|
|
except Exception:
|
|
|
|
traceback.print_exc()
|
|
|
|
return
|
|
|
|
|
|
|
|
if msg.dst != messages.ROOT_ADDRESS and msg.dst != messages.PARENT_ADDRESS:
|
|
|
|
print('Received message for forwarding:', msg)
|
|
|
|
# todo: this message isn't for us, forward it
|
|
|
|
return
|
|
|
|
|
2022-04-04 14:48:43 +02:00
|
|
|
print('Received message:', msg)
|
2022-04-15 20:57:11 +02:00
|
|
|
|
2023-10-03 17:23:29 +02:00
|
|
|
src_node, created = MeshNode.objects.get_or_create(address=msg.src)
|
|
|
|
|
2022-04-04 14:48:43 +02:00
|
|
|
if isinstance(msg, messages.MeshSigninMessage):
|
2023-10-03 17:23:29 +02:00
|
|
|
self.uplink_node = src_node
|
|
|
|
# log message, since we will not log it further down
|
|
|
|
self.log_received_message(src_node, msg)
|
|
|
|
|
|
|
|
# inform signed in uplink node about its layer
|
|
|
|
self.send_msg(messages.MeshLayerAnnounceMessage(
|
2022-04-15 20:02:42 +02:00
|
|
|
src=messages.ROOT_ADDRESS,
|
2022-04-06 17:25:46 +02:00
|
|
|
dst=msg.src,
|
|
|
|
layer=messages.NO_LAYER
|
2022-04-06 22:56:08 +02:00
|
|
|
))
|
2022-04-15 20:57:11 +02:00
|
|
|
|
2023-10-03 17:23:29 +02:00
|
|
|
# add signed in uplink node to broadcast route
|
|
|
|
async_to_sync(self.channel_layer.group_add)('mesh_broadcast', self.channel_name)
|
2022-04-15 20:57:11 +02:00
|
|
|
|
2023-10-03 17:23:29 +02:00
|
|
|
# add this node as a destination that this uplink handles (duh)
|
|
|
|
self.add_dst_nodes((src_node.address, ))
|
2022-04-15 20:57:11 +02:00
|
|
|
|
2023-10-02 22:02:25 +02:00
|
|
|
return
|
|
|
|
|
2023-10-03 17:23:29 +02:00
|
|
|
if self.uplink_node is None:
|
|
|
|
print('Expected sign-in message, but got a different one!')
|
|
|
|
self.close()
|
|
|
|
return
|
2023-10-02 22:02:25 +02:00
|
|
|
|
2023-10-03 17:23:29 +02:00
|
|
|
self.log_received_message(src_node, msg)
|
2022-04-15 20:02:42 +02:00
|
|
|
|
2023-10-03 17:23:29 +02:00
|
|
|
def uplink_change(self, data):
|
|
|
|
# message handler: if we are not the given uplink, leave this group
|
|
|
|
if data["uplink"] != self.uplink_node.address:
|
|
|
|
group = self.group_name_for_node(data["address"])
|
|
|
|
print('leaving uplink group...')
|
|
|
|
async_to_sync(self.channel_layer.group_discard)(group, self.channel_name)
|
|
|
|
|
|
|
|
def log_received_message(self, src_node: MeshNode, msg: messages.Message):
|
2022-04-15 20:57:11 +02:00
|
|
|
NodeMessage.objects.create(
|
2023-10-03 17:23:29 +02:00
|
|
|
uplink_node=self.uplink_node,
|
|
|
|
src_node=src_node,
|
2022-04-15 20:02:42 +02:00
|
|
|
message_type=msg.msg_id,
|
|
|
|
data=msg.tojson()
|
|
|
|
)
|
2022-04-15 20:57:11 +02:00
|
|
|
|
2023-10-03 17:23:29 +02:00
|
|
|
def add_dst_nodes(self, addresses):
|
|
|
|
# add ourselves to this one
|
|
|
|
for address in addresses:
|
|
|
|
# create group name for this address
|
|
|
|
group = self.group_name_for_node(address)
|
|
|
|
|
|
|
|
# if we aren't handling this address yet, join the group
|
|
|
|
if address not in self.dst_nodes:
|
|
|
|
async_to_sync(self.channel_layer.group_add)(group, self.channel_name)
|
|
|
|
self.dst_nodes.add(address)
|
|
|
|
|
|
|
|
# tell other consumers to leave the group
|
|
|
|
async_to_sync(self.channel_layer.group_send)(group, {
|
|
|
|
"type": "uplink_change",
|
|
|
|
"node": address,
|
|
|
|
"uplink": self.uplink_node.address
|
|
|
|
})
|
|
|
|
|
|
|
|
# tell the node to dump its current information
|
|
|
|
self.send_msg(
|
|
|
|
messages.ConfigDumpMessage(
|
|
|
|
src=messages.ROOT_ADDRESS,
|
|
|
|
dst=address,
|
|
|
|
)
|
|
|
|
)
|
|
|
|
|
|
|
|
# add the stuff to the db as well
|
|
|
|
MeshNode.objects.filter(address__in=addresses).update(route_id=self.uplink_node.address)
|
|
|
|
|
|
|
|
def group_name_for_node(self, address):
|
|
|
|
return 'mesh_%s' % address.replace(':', '-')
|
2022-04-15 20:57:11 +02:00
|
|
|
|
2022-04-15 21:06:57 +02:00
|
|
|
def remove_route(self, route_address):
|
|
|
|
MeshNode.objects.filter(route_id=route_address).update(route_id=None)
|