import socket import threading import queue import random import struct import time import math import pyproj import vdes_protocol class vdes_server: def __init__(self, host='localhost', port=10): self.host = host self.port = port self.bufsize = 32768 self.addr = (self.host, self.port) self.tcp_send_queue = queue.Queue(32) self.tcp_recv_queue = queue.Queue(16) self.tcp_recv_buffer = bytearray() self.tcp_recv_pointer = 0 self.rf_recv_queue = queue.Queue(32) self.rf_recv_buffer = bytearray() self.rf_recv_pointer = 0 self.tcpCliSock = None self.tcpSerSock = socket.socket(socket.AF_INET, socket.SOCK_STREAM) self.tcpSerSock.bind(self.addr) self.tcpSerSock.listen(2) self.my_lon = 121.167 self.my_lat = 31.187 self.my_heading = 270 self.gps_data = vdes_protocol.pkt_gps(self.my_lon, self.my_lat, 0, 0, self.my_heading) # 模拟 40 条船只以本设备为中心 4 ~ 7 海里距离,角度均分,方向随机为离心或向心 self.ais_data = list() for i in range(40): radius = random.randint(40, 70)/10*1852 # distributed in 4 ~ 7 nautical miles ship_direction = random.randint(0, 1) ship_heading = (i*9 + ship_direction*180) % 360 # Use WGS84 ellipsoid. g = pyproj.Geod(ellps='WGS84') # Forward transformation - returns longitude, latitude, back azimuth of terminus points ship_lon, ship_lat, az = g.fwd(self.my_lon, self.my_lat, ship_heading, radius) ship_mmsi = random.randint(100000000, 999999999) ship_speed = random.randint(14, 30) self.ais_data.append(vdes_protocol.pkt_ais_msg_1(ship_mmsi, ship_speed, ship_lon, ship_lat, ship_heading)) def recv_parse(self): pass def send(self, stop_event): print("send started.") while not stop_event.is_set(): if self.tcp_send_queue.empty(): time.sleep(0.01) continue data = self.tcp_send_queue.get() if self.tcpCliSock: self.tcpCliSock.send(data) def make_response(self, stop_event): print("make_response started.") while not stop_event.is_set(): if self.tcp_recv_queue.empty(): time.sleep(0.1) continue recv_data += self.tcp_recv_queue.get() def make_gps(self, stop_event): print("make_gps started.") while not stop_event.is_set(): self.tcp_send_queue.put(self.gps_data.get_packet()) self.tcp_send_queue.task_done() time.sleep(1) def make_ais(self, stop_event): print("make_ais started.") while not stop_event.is_set(): for data in self.ais_data: if stop_event.is_set(): break data = data.get_packet(data.sog * 1852 / 3600 * 2) # 速度放大10倍 self.tcp_send_queue.put(data) self.tcp_send_queue.task_done() time.sleep(0.05) def start(self): threads = [] while True: print('Waiting for connection...') self.tcpCliSock, addr = self.tcpSerSock.accept() print('...connected from:', addr) stop_event_response = threading.Event() threads.append(threading.Thread(target=self.make_response, args=(stop_event_response,))) stop_event_gps = threading.Event() threads.append(threading.Thread(target=self.make_gps, args=(stop_event_gps,))) stop_event_ais = threading.Event() threads.append(threading.Thread(target=self.make_ais, args=(stop_event_ais,))) stop_event_send = threading.Event() threads.append(threading.Thread(target=self.send, args=(stop_event_send,))) for t in threads: t.start() try: while True: recv_data = self.tcpCliSock.recv(self.bufsize) if not recv_data: print("Client disconnected") break self.tcp_recv_queue.put(recv_data) except (ConnectionResetError, BrokenPipeError): print("Client forcibly disconnected") except Exception as e: print(f"Error: {e}") finally: self.tcpCliSock.close() stop_event_response.set() stop_event_gps.set() stop_event_ais.set() stop_event_send.set() for t in threads: t.join() threads.clear() self.tcpCliSock = None if __name__ == "__main__": server = vdes_server() try: server.start() except KeyboardInterrupt: print("Server stopped by user") finally: if server.tcpCliSock: server.tcpCliSock.close() server.tcpSerSock.close()