141 lines
5.0 KiB
Python
141 lines
5.0 KiB
Python
|
|
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()
|