P2P聊天中get_clients_from_signal_server无数据且程序冻结问题
P2P聊天获取客户端列表时程序冻结无响应问题
我在P2P聊天中调用get_clients_from_signal_server时无法获取任何数据;当接收已连接用户列表时,程序直接冻结,服务器停止响应,终端仅显示源码但无消息。
信号服务器(signal_server)代码
import socket import threading import json clients = {} lock = threading.Lock() def handle_client(data, addr, sock): global clients try: message = json.loads(data.decode('utf-8')) print(f"Received message from {addr}: {message}") # Debugging statement if 'register' in message: with lock: clients[message['name']] = addr response = {'status': 'registered', 'clients': list(clients.keys())} sock.sendto(json.dumps(response).encode('utf-8'), addr) print(f"Registered {message['name']} from {addr}") # Debugging statement elif 'get_clients' in message: response = {'status': 'registered', 'clients': list(clients.keys())} sock.sendto(json.dumps(response).encode('utf-8'), addr) print(f"Sent client list to {addr}") # Debugging statement except Exception as e: print(f"Error handling client {addr}: {e}") def signal_server(host, port): server_socket = socket.socket(socket.AF_INET, socket.SOCK_DGRAM) server_socket.bind((host, port)) print(f"Signal server started on {host}:{port}") while True: try: data, addr = server_socket.recvfrom(1024) threading.Thread(target=handle_client, args=(data, addr, server_socket)).start() except Exception as e: print(f"Error in main loop: {e}") if __name__ == '__main__': signal_host = 'localhost' signal_port = 8001 signal_server(signal_host, signal_port)
P2P客户端(P2PClient)代码
import argparse import socket import sys import time import json import threading shutdown = False class Message: def __init__(self, **data): self.status = 'online' for param, value in data.items(): setattr(self, param, value) self.curr_time = time.strftime("%Y-%m-%d-%H.%M.%S", time.localtime()) def to_json(self): return json.dumps(self, default=lambda o: o.__dict__, sort_keys=True, indent=4) class P2PClient: def __init__(self, host, port, name=None, signal_server_address=None): self.current_connection = None self.client_address = (host, port) self.name = name if name else f"{host}:{port}" self.signal_server_address = signal_server_address self.socket = socket.socket(socket.AF_INET, socket.SOCK_DGRAM) self.socket.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1) self.socket.bind(self.client_address) def register_with_signal_server(self): if self.signal_server_address: message = {"register": True, "name": self.name} self.socket.sendto(json.dumps(message).encode('utf-8'), self.signal_server_address) data, _ = self.socket.recvfrom(1024) response = json.loads(data.decode('utf-8')) if response['status'] == 'registered': print(f"Registered with signal server. Current clients: {response['clients']}") def get_clients_from_signal_server(self): if self.signal_server_address: message = {"get_clients": True} print(f"Requesting client list from signal server at {self.signal_server_address}") # Debugging statement self.socket.sendto(json.dumps(message).encode('utf-8'), self.signal_server_address) data, _ = self.socket.recvfrom(1024) response = json.loads(data.decode('utf-8')) print(f"Received clients from signal server: {response['clients']}") return response['clients'] return [] ... def run(self): global shutdown self.register_with_signal_server() recv_thread = threading.Thread(target=self.receive) recv_thread.start() while not shutdown: print("1. Get list of clients") print("2. Connect to client") print("3. Exit") choice = input("Choose an option: ") if choice == "1": clients = self.get_clients_from_signal_server() print("Available clients:") for client_name in clients: print(f" - {client_name}") elif choice == "2": self.connect() elif choice == "3": shutdown = True print("Exiting...") else: print("Invalid option. Please try again.") time.sleep(1) recv_thread.join() if __name__ == '__main__': parser = argparse.ArgumentParser() parser.add_argument("-ho", "--host", help="p2p client host ip address, like 127.0.0.1") parser.add_argument("-p", "--port", help="p2p client host port, like 8001") parser.add_argument("-sh", "--signalhost", help="signal server host ip address, like 127.0.0.1") parser.add_argument("-sp", "--signalport", help="signal server port, like 8000") args = parser.parse_args() try: host = args.host port = int(args.port) signal_host = args.signalhost signal_port = int(args.signalport) name = input("Name: ").strip() signal_server_address = (signal_host, signal_port) p2p_client = P2PClient(host, port, name=name, signal_server_address=signal_server_address) p2p_client.run() except (TypeError, ValueError): print("Incorrect arguments values, use --help/-h for more info.")
客户端启动脚本(Bat)
python p2p_client.py --host localhost --port 8002 --signalhost localhost --signalport 8001
python p2p_client.py --host localhost --port 8003 --signalhost localhost --signalport 8001
问题原因及修复方案
核心问题
- UDP Socket接收冲突:客户端主线程调用
get_clients_from_signal_server时,会阻塞在recvfrom上,同时接收线程也在使用同一个Socket调用recvfrom。信号服务器的响应可能被接收线程抢先读取,导致主线程无限等待,程序冻结。 - 无超时机制:UDP本身不保证可靠传输,若响应丢失,客户端会一直阻塞在
recvfrom,没有超时退出逻辑。
修复步骤
1. 给Socket添加超时并捕获异常
在P2PClient的__init__方法中设置Socket超时:
self.socket.settimeout(5) # 设置5秒超时,可按需调整
修改get_clients_from_signal_server方法,捕获超时和其他异常:
def get_clients_from_signal_server(self): if self.signal_server_address: message = {"get_clients": True} print(f"Requesting client list from signal server at {self.signal_server_address}") try: self.socket.sendto(json.dumps(message).encode('utf-8'), self.signal_server_address) data, _ = self.socket.recvfrom(1024) response = json.loads(data.decode('utf-8')) print(f"Received clients from signal server: {response['clients']}") return response['clients'] except socket.timeout: print("超时:未收到信号服务器响应") return [] except Exception as e: print(f"获取客户端列表失败:{e}") return [] return []
2. 分离信号通信与P2P通信Socket
使用两个独立Socket分别处理信号服务器交互和P2P消息接收,避免接收冲突:
class P2PClient: def __init__(self, host, port, name=None, signal_server_address=None): self.current_connection = None self.client_address = (host, port) self.name = name if name else f"{host}:{port}" self.signal_server_address = signal_server_address # P2P通信Socket(绑定端口) self.p2p_socket = socket.socket(socket.AF_INET, socket.SOCK_DGRAM) self.p2p_socket.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1) self.p2p_socket.bind(self.client_address) # 信号服务器通信Socket(不绑定端口,系统自动分配) self.signal_socket = socket.socket(socket.AF_INET, socket.SOCK_DGRAM) self.signal_socket.settimeout(5)
修改register_with_signal_server和get_clients_from_signal_server使用signal_socket,接收线程使用p2p_socket:
def register_with_signal_server(self): if self.signal_server_address: message = {"register": True, "name": self.name} try: self.signal_socket.sendto(json.dumps(message).encode('utf-8'), self.signal_server_address) data, _ = self.signal_socket.recvfrom(1024) response = json.loads(data.decode('utf-8')) if response['status'] == 'registered': print(f"已注册到信号服务器,当前客户端:{response['clients']}") except socket.timeout: print("注册超时:未收到信号服务器响应") except Exception as e: print(f"注册失败:{e}") # 接收线程逻辑 def receive(self): global shutdown while not shutdown: try: data, addr = self.p2p_socket.recvfrom(1024) print(f"收到来自{addr}的消息:{data.decode('utf-8')}") except socket.timeout: continue except Exception as e: if not shutdown: print(f"接收消息失败:{e}")
3. 添加请求重传机制(可选)
针对UDP不可靠性,添加重传逻辑提升成功率:
def get_clients_from_signal_server(self): if self.signal_server_address: message = {"get_clients": True} print(f"请求客户端列表:{self.signal_server_address}") retries = 2 # 最多重传2次 for i in range(retries + 1): try: self.signal_socket.sendto(json.dumps(message).encode('utf-8'), self.signal_server_address) data, _ = self.signal_socket.recvfrom(1024) response = json.loads(data.decode('utf-8')) print(f"收到客户端列表:{response['clients']}") return response['clients'] except socket.timeout: if i < retries: print(f"超时,重试({i+1}/{retries})...") else: print("已达最大重试次数,未收到响应") except Exception as e: print(f"获取列表失败:{e}") return [] return [] return []
内容的提问来源于stack exchange,提问作者Павел Бархотов
相关产品推荐
相关产品推荐

