You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

问题原因及修复方案

核心问题

  1. UDP Socket接收冲突:客户端主线程调用get_clients_from_signal_server时,会阻塞在recvfrom上,同时接收线程也在使用同一个Socket调用recvfrom。信号服务器的响应可能被接收线程抢先读取,导致主线程无限等待,程序冻结。
  2. 无超时机制: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,提问作者Павел Бархотов

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.21 21:22:05