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()