240 lines
8.9 KiB
Python
240 lines
8.9 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):
|
||
|
|
global frame_index
|
||
|
|
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:
|
||
|
|
if ais_frames:
|
||
|
|
# print(ais_frames, type(ais_frames))
|
||
|
|
frame = ais_frames[frame_index]
|
||
|
|
try:
|
||
|
|
conn.sendall(frame)
|
||
|
|
print(f"发送报文帧[{frame_index}]:", frame.hex())
|
||
|
|
except (ConnectionResetError, ConnectionAbortedError, BrokenPipeError):
|
||
|
|
print("发送失败,客户端断开")
|
||
|
|
break
|
||
|
|
|
||
|
|
# 轮流取下一帧
|
||
|
|
frame_index = (frame_index + 1) % len(ais_frames)
|
||
|
|
time.sleep(5)
|
||
|
|
|
||
|
|
# pass
|
||
|
|
|
||
|
|
|
||
|
|
|
||
|
|
# 生成并发送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 = 'localhost'
|
||
|
|
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("服务端停止")
|
||
|
|
def load_frames(filepath):
|
||
|
|
with open(filepath, "rb") as f:
|
||
|
|
data = f.read()
|
||
|
|
|
||
|
|
delimiter = b"####" # 帧起始标记
|
||
|
|
frames = data.split(delimiter)
|
||
|
|
|
||
|
|
# 加回起始符,过滤掉空的
|
||
|
|
frames = [delimiter + frame for frame in frames if frame]
|
||
|
|
return frames
|
||
|
|
|
||
|
|
if __name__ == "__main__":
|
||
|
|
# 全局变量:预加载所有帧
|
||
|
|
ais_frames = load_frames("./tcp_data/TCP_GNSS消息.bin")
|
||
|
|
frame_index = 0 # 当前取到第几帧
|
||
|
|
send_data_continuously()
|