wu_xinjun
2022-06-21 ee0fccfdd6efc0cdc1cc81dbc67a2bfaae8c7ff1
optimize for routing key
6个文件已修改
115 ■■■■ 已修改文件
scripts_launch/TEST_startClient.bat 5 ●●●●● 补丁 | 查看 | 原始文档 | blame | 历史
scripts_launch/TEST_startServer.bat 8 ●●●● 补丁 | 查看 | 原始文档 | blame | 历史
src/Server.py 22 ●●●● 补丁 | 查看 | 原始文档 | blame | 历史
src/engine/AcquisitionMain.py 18 ●●●● 补丁 | 查看 | 原始文档 | blame | 历史
src/engine/MessageThreads.py 2 ●●● 补丁 | 查看 | 原始文档 | blame | 历史
src/engine/devices.py 60 ●●●●● 补丁 | 查看 | 原始文档 | blame | 历史
scripts_launch/TEST_startClient.bat
@@ -16,11 +16,12 @@
echo change to path: %cd%
set MQhost=localhost
@REM set MQhost=localhost
set MQhost=47.92.33.19
set MQuser=chai
set MQpassword=password123
set QUEUEname=aerial_rpc
set DEVICEname=client.client.client
set DEVICEname=bunker0.block0.device0
set SERVERhost=http://192.168.1.224/api
set MICROWAVEhost=localhost
set LOGfile=ClientTEST.log
scripts_launch/TEST_startServer.bat
@@ -16,14 +16,14 @@
echo change to path: %cd%
set MQhost=localhost
@REM set MQhost=47.92.33.19
@REM set MQhost=localhost
set MQhost=47.92.33.19
set MQuser=chai
set MQpassword=password123
set WEBhost=0.0.0.0
set WEBport=5000
set SAVEfolder=./data/test/
set SERVERname=testserver
set SAVEfolder=./data/server0/
set SERVERname=server0
set LOGfile=ServerTEST.log
@REM Launch the Server
src/Server.py
@@ -70,12 +70,12 @@
                                    MQparams=self.parameters,
                                    exchange_name='logs_topic',
                                    exchange_type='topic',
                                    MQqueue_name=f'logs_queue',
                                    binding_keys=['server.logs.test0'],
                                    MQqueue_name=f'{self.server_name}_logs_queue',
                                    binding_keys=[f'{self.server_name}.logs.test0'],
                                    message_queue=self.LogsReceiver_queue
                                    )
        self.LogsReceiver_Thread.setDaemon(True)
        self.LogsReceiver_Thread.setName("logs_thread")
        self.LogsReceiver_Thread.setName("Thread-logsReceiver")
        self.LogsReceiver_Thread.start()
        # launch images receiver thread
@@ -83,12 +83,12 @@
                                    MQparams=self.parameters,
                                    exchange_name='images_topic',
                                    exchange_type='topic',
                                    MQqueue_name=f'images_queue',
                                    binding_keys=["server.images.test0"],
                                    MQqueue_name=f'{self.server_name}_images_queue',
                                    binding_keys=[f"{self.server_name}.images.test0"],
                                    message_queue=self.ImagesReceiver_queue
                                    )
        self.ImagesReceiver_Thread.setDaemon(True)
        self.ImagesReceiver_Thread.setName("images_thread")
        self.ImagesReceiver_Thread.setName("Thread-imagesReceiver")
        self.ImagesReceiver_Thread.start()
        # launch plys receiver thread
@@ -96,12 +96,12 @@
                                    MQparams=self.parameters,
                                    exchange_name='plys_topic',
                                    exchange_type='topic',
                                    MQqueue_name=f'plys_queue',
                                    binding_keys=['server.plys.test0'],
                                    MQqueue_name=f'{self.server_name}_plys_queue',
                                    binding_keys=[f'{self.server_name}.plys.test0'],
                                    message_queue=self.PlysReceiver_queue
                                    )
        self.PlysReceiver_Thread.setDaemon(True)
        self.PlysReceiver_Thread.setName("plys_thread")
        self.PlysReceiver_Thread.setName("Thread-plysReceiver")
        self.PlysReceiver_Thread.start()
        # launch decode thread
@@ -114,7 +114,7 @@
                                    )
        self.MsgDecoder_Thread.setDaemon(True)
        self.MsgDecoder_Thread.setName("decode_thread")
        self.MsgDecoder_Thread.setName("Thread-dataDecoder")
        self.MsgDecoder_Thread.start()
    
@@ -128,7 +128,7 @@
                                    message_queue=self.CommandsSender_queue
                        )
        self.CommandsSender_Thread.setDaemon(True)
        self.CommandsSender_Thread.setName("commands_thread")
        self.CommandsSender_Thread.setName("Thread-commandsSender")
        self.CommandsSender_Thread.start()
src/engine/AcquisitionMain.py
@@ -133,12 +133,12 @@
                                    MQparams=self.parameters,
                                    exchange_name='commands_topic',
                                    exchange_type='topic',
                                    MQqueue_name="commands_queue",
                                    MQqueue_name="server0_commands_queue",
                                    binding_keys=[f"{self.device_id}"],
                                    message_queue=self.CommandsReceiver_queue
                                    )
        self.CommandsReceiver_Thread.setDaemon(True)
        self.CommandsReceiver_Thread.setName("commands_thread")
        self.CommandsReceiver_Thread.setName("Thread-commandsReceiver")
        self.CommandsReceiver_Thread.start()
        self.threads.append(self.CommandsReceiver_Thread)
@@ -153,7 +153,7 @@
                                message_queue=self.LogsSender_queue
                        )
        self.LogsSender_Thread.setDaemon(True)
        self.LogsSender_Thread.setName("logs_thread")
        self.LogsSender_Thread.setName("Thread-logsSender")
        self.LogsSender_Thread.start()
        self.threads.append(self.LogsSender_Thread)
@@ -167,7 +167,7 @@
                                message_queue=self.ImagesSender_queue
                        )
        self.ImagesSender_Thread.setDaemon(True)
        self.ImagesSender_Thread.setName("image_thread")
        self.ImagesSender_Thread.setName("Thread-imageSender")
        self.ImagesSender_Thread.start()
        self.threads.append(self.ImagesSender_Thread)
@@ -180,7 +180,7 @@
                                message_queue=self.PlysSender_queue
                        )
        self.PlysSender_Thread.setDaemon(True)
        self.PlysSender_Thread.setName("plys_thread")
        self.PlysSender_Thread.setName("Thread-plysSender")
        self.PlysSender_Thread.start()
        self.threads.append(self.PlysSender_Thread)
@@ -196,8 +196,8 @@
                        device_id=self.device_id,
                        logs_queue=self.LogsSender_queue,
                        images_queue=self.ImagesSender_queue,
                        init_binding_key="client.client.client",
                        init_routing_key="server"
                        init_binding_key=self.device_id,
                        init_routing_key="server0"
        )
        self.dc = dustCleanerController(
@@ -206,7 +206,7 @@
                    timeout=self.timeout,
                    logs_queue=self.LogsSender_queue,
                    device_id=self.device_id,
                    routing_key="server"
                    routing_key="server0"
                    )
        self.mc = microwaveController(
@@ -217,7 +217,7 @@
                    device_id=self.device_id,
                    logs_queue=self.LogsSender_queue,
                    plys_queue=self.PlysSender_queue,
                    routing_key="server"
                    routing_key="server0"
        )
        
src/engine/MessageThreads.py
@@ -82,7 +82,7 @@
                        )
                        timestamp = datetime.now().strftime('%Y-%m-%d-%H-%M-%S')
                        caption = f"[x] {timestamp} | Sending message to[{msg_routing_key}]"
                        caption = f"[x] {timestamp} | Sending message to [{msg_routing_key}]"
                        print(caption)
                    else:
src/engine/devices.py
@@ -204,63 +204,3 @@
                        logs_queue=self.logs_queue,
                        source = self.device_id,
                        routing_key=self.routing_key+".logs.test0")
# class CaptureThread(threading.Thread):
#     """
#     camera capture thread
#     """
#     def __init__(self, cameras, msg_queue, interval, device_id,image_quality = 95):
#         """init the thread"""
#         threading.Thread.__init__(self)
#         self.cameras = cameras
#         self.msg_queue = msg_queue
#         self.interval = interval
#         self.image_quality = image_quality
#         self.is_stop = False
#         self.device_id = device_id
#     def stop(self):
#         """stop the thread"""
#         self.is_stop = True
# # TODO根据新的数据结构修改队消息队列。
#     def run(self):
#         print("CaptureThread started")
#         encode_param = [int(cv2.IMWRITE_JPEG_QUALITY), self.image_quality]
#         self.next_time = time.time() + self.interval
#         while not self.is_stop:
#             this_time = time.time()
#             if this_time > self.next_time:
#                 for i in range(len(self.cameras)):
#                     _, frame = self.cameras[i].read()
#                     capture_time = time.time()
#                     _, encimg = cv2.imencode('.jpg', frame, encode_param)
#                     if frame is not None:
#                         np_array = np.asarray(encimg)  # frame 转换为numpy数组
#                         byte_array = np_array.tobytes()  # numpy数组转换为byte数组
#                         # convert the image to pb.DataField
#                         imagedata = pb.DataField()
#                         imagedata.data = byte_array
#                         imagedata.length = len(byte_array)
#                         data_str = imagedata.SerializeToString()
#                         # create protobuf message
#                         data_msg = pb.CapturedData()
#                         data_msg.generate_time = capture_time
#                         data_msg.generate_device = self.device_id
#                         data_msg.camera_port = i
#                         data_msg.image_data = data_str
#                         # add message to queue
#                         self.msg_queue.put(data_msg)
#                 self.next_time += self.interval
#             time.sleep(0.1)
#         for camera in self.cameras:
#             camera.release()
#         print("CaptureThread stopped")