209 lines
7.8 KiB
Python
209 lines
7.8 KiB
Python
import struct
|
|
import socket
|
|
import time
|
|
import random
|
|
import math
|
|
import threading
|
|
from vdes_message_unpack import VdesMessageUnpack
|
|
|
|
# ------------------------ 全局队列和锁 ------------------------
|
|
vdes_message_queue = []
|
|
message_queue_lock = threading.Lock()
|
|
|
|
# ------------------------ 船只初始化 ------------------------
|
|
SHIP_MMSIS = [random.randint(336295939, 336295969) for _ in range(10)]
|
|
SHIP_STATUS = {}
|
|
SHIP_LIST = [
|
|
(336295939 + i, round(random.uniform(30.0, 34.0), 6),
|
|
round(random.uniform(123.5, 126.0), 6), random.randint(0, 359))
|
|
for i in range(30)
|
|
]
|
|
for ship in SHIP_LIST:
|
|
mmsi, init_lat, init_lon, heading = ship
|
|
SHIP_STATUS[mmsi] = {
|
|
'current_lat': init_lat,
|
|
'current_lon': init_lon,
|
|
'heading': heading,
|
|
'speed': 280.24 # 加速100倍用于演示
|
|
}
|
|
|
|
# ------------------------ 当前坐标 ------------------------
|
|
YANGTZE_ESTUARY_LAT = 31.23
|
|
YANGTZE_ESTUARY_LON = 123.87
|
|
current_lat = YANGTZE_ESTUARY_LAT
|
|
current_lon = YANGTZE_ESTUARY_LON
|
|
own_heading = 90 # 固定航向
|
|
|
|
# ------------------------ 工具函数 ------------------------
|
|
def generate_random_position():
|
|
global current_lat, current_lon, own_heading
|
|
speed_knots = 300.24
|
|
speed_kmh = speed_knots * 1.852
|
|
speed_kmps = speed_kmh / 3600
|
|
distance_km = speed_kmps * 1
|
|
heading_rad = math.radians(own_heading)
|
|
delta_lat = distance_km / 111.0 * math.cos(heading_rad)
|
|
delta_lon = distance_km / (111.0 * math.cos(math.radians(current_lat))) * math.sin(heading_rad)
|
|
current_lat += delta_lat
|
|
current_lon += delta_lon
|
|
return int(current_lon * 10000000), int(current_lat * 10000000)
|
|
|
|
def generate_ship_position(mmsi):
|
|
ship = SHIP_STATUS[mmsi]
|
|
speed_kmh = ship['speed'] * 1.852
|
|
distance_km = speed_kmh / 3600
|
|
heading_rad = math.radians(ship['heading'])
|
|
delta_lat = distance_km / 111.0 * math.cos(heading_rad)
|
|
delta_lon = distance_km / (111.0 * math.cos(math.radians(ship['current_lat']))) * math.sin(heading_rad)
|
|
ship['current_lat'] += delta_lat
|
|
ship['current_lon'] += delta_lon
|
|
return ship['current_lon'], ship['current_lat']
|
|
|
|
def create_tcp_packet(packet_type, content):
|
|
header = b'####'
|
|
type_field = struct.pack('>I', packet_type)
|
|
length_field = struct.pack('>I', len(content))
|
|
return header + type_field + length_field + content
|
|
|
|
def create_gps_packet():
|
|
lon, lat = generate_random_position()
|
|
utc_time = int(time.time())
|
|
content = struct.pack('>QiiiHHHBBIII',
|
|
utc_time, lon, lat, 50, 10, 8, 7, 1, 3, 2550, 1024, own_heading)
|
|
return content
|
|
|
|
def create_ais_vdes_info_packet(demod_content, Demodulator_id):
|
|
snr = 25.75
|
|
snr_int = int(snr)
|
|
snr_frac = int((snr - snr_int) * 100)
|
|
snr_packed = (snr_int & 0xFF) << 8 | (snr_frac & 0xFF)
|
|
content = struct.pack('>BBBHHH', 12, 1, Demodulator_id, 1234, 5678, snr_packed)
|
|
return content + demod_content
|
|
|
|
def encode_ais_demod_fields(SHIP_MMSI, lat, lon, heading):
|
|
bits = []
|
|
bits.append(format(1, '06b'))
|
|
bits.append(format(0, '02b'))
|
|
bits.append(format(SHIP_MMSI, '030b'))
|
|
bits.append(format(8, '04b'))
|
|
bits.append(format(128 & 0xFF, '08b'))
|
|
bits.append(format(102, '010b'))
|
|
bits.append(str(1))
|
|
bits.append(format(int(lon * 60 * 10000) & 0x0FFFFFFF, '028b'))
|
|
bits.append(format(int(lat * 60 * 10000) & 0x07FFFFFF, '027b'))
|
|
bits.append(format(min(heading, 3600), '012b'))
|
|
bits.append(format(min(heading, 511), '09b'))
|
|
bits.append(format(30, '06b'))
|
|
bits.append(format(0, '02b'))
|
|
bits.append(format(0, '03b'))
|
|
bits.append(str(0))
|
|
bits.append(format(0, '019b'))
|
|
bit_string = ''.join(bits)
|
|
binary_int = int(bit_string, 2)
|
|
byte_length = (len(bit_string) + 7) // 8
|
|
return binary_int.to_bytes(byte_length, 'big')
|
|
|
|
def encode_vdes_demod_fields(Type, Fragment_num, MMSI, Payload):
|
|
source_id = 12345678
|
|
satellite_id = 1
|
|
session_id = 100
|
|
dest_id = MMSI
|
|
fragment_num = Fragment_num
|
|
binary_payload = Payload
|
|
payload_size = 4 + 1 + 1 + 4 + 2 + len(binary_payload)
|
|
content = struct.pack('>B H I B B I H',
|
|
Type, payload_size, source_id, satellite_id, session_id, dest_id, fragment_num)
|
|
return content + binary_payload
|
|
|
|
# ------------------------ 线程函数 ------------------------
|
|
def send_messages(conn):
|
|
try:
|
|
while True:
|
|
with message_queue_lock:
|
|
if vdes_message_queue:
|
|
packet = vdes_message_queue.pop(0)
|
|
else:
|
|
packet = None
|
|
if packet:
|
|
try:
|
|
conn.sendall(packet)
|
|
except (ConnectionResetError, ConnectionAbortedError, BrokenPipeError):
|
|
print("发送失败,客户端断开")
|
|
break
|
|
else:
|
|
time.sleep(1)
|
|
# 生成并发送AIS报文
|
|
# selected_ship = random.choice(SHIP_LIST)
|
|
# mmsi = selected_ship[0]
|
|
# lon, lat = generate_ship_position(mmsi)
|
|
# heading = SHIP_STATUS[mmsi]['heading']
|
|
# ais_demod_content = encode_ais_demod_fields(mmsi, lat, lon, heading)
|
|
# ais_data = create_ais_vdes_info_packet(ais_demod_content, 4)
|
|
# ais_packet = create_tcp_packet(12, ais_data)
|
|
# #生成并发送GPS报文
|
|
# gps_content = create_gps_packet()
|
|
# gps_packet = create_tcp_packet(7, gps_content)
|
|
# try:
|
|
# conn.sendall(ais_packet)
|
|
# conn.sendall(gps_packet)
|
|
# except (ConnectionResetError, ConnectionAbortedError, BrokenPipeError):
|
|
# print("发送失败,客户端断开")
|
|
# break
|
|
time.sleep(0.5)
|
|
finally:
|
|
print("发送线程退出")
|
|
|
|
def handle_client(conn, addr):
|
|
print(f"客户端连接: {addr}")
|
|
unpacker = VdesMessageUnpack()
|
|
fragment_num = 0
|
|
send_thread = threading.Thread(target=send_messages, args=(conn,))
|
|
send_thread.daemon = True
|
|
send_thread.start()
|
|
try:
|
|
while True:
|
|
try:
|
|
data = conn.recv(65536)
|
|
if not data:
|
|
print("客户端关闭连接")
|
|
break
|
|
# 解析报文并加入发送队列
|
|
frame_type, message = unpacker.unpack_tcp_packet(data)
|
|
mmsi = unpacker.get_mmsi()
|
|
Type = 30 if frame_type == 0 else 31 if frame_type == 1 else 32
|
|
vdes_demod_content = encode_vdes_demod_fields(Type, fragment_num, mmsi, message)
|
|
vdes_data = create_ais_vdes_info_packet(vdes_demod_content, 0)
|
|
vdes_packet = create_tcp_packet(12, vdes_data)
|
|
fragment_num += 1
|
|
with message_queue_lock:
|
|
vdes_message_queue.append(vdes_packet)
|
|
except (ConnectionResetError, ConnectionAbortedError):
|
|
print("客户端断开连接")
|
|
break
|
|
except Exception as e:
|
|
print(f"解析报文错误: {e}")
|
|
finally:
|
|
conn.close()
|
|
print("客户端连接关闭")
|
|
|
|
# ------------------------ 服务端启动 ------------------------
|
|
def send_data_continuously():
|
|
host = '127.0.0.1'
|
|
port = 10
|
|
with socket.socket(socket.AF_INET, socket.SOCK_STREAM) as s:
|
|
s.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
|
|
s.bind((host, port))
|
|
s.listen(5)
|
|
print(f"服务端启动 {host}:{port}")
|
|
try:
|
|
while True:
|
|
conn, addr = s.accept()
|
|
client_thread = threading.Thread(target=handle_client, args=(conn, addr))
|
|
client_thread.daemon = True
|
|
client_thread.start()
|
|
except KeyboardInterrupt:
|
|
print("服务端停止")
|
|
|
|
if __name__ == "__main__":
|
|
send_data_continuously()
|