修复多连接的问题
This commit is contained in:
@@ -525,40 +525,69 @@ def udp_echo_server_test(args):
|
||||
|
||||
# ==================== 压力测试 ====================
|
||||
|
||||
def stress_test_thread(client_id, args, results):
|
||||
def stress_test_thread(client_id, args, results, barrier):
|
||||
"""压力测试线程"""
|
||||
thread_name = f"Client-{client_id}"
|
||||
sock = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
|
||||
sock.settimeout(10)
|
||||
sock.settimeout(5)
|
||||
sock.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
|
||||
|
||||
try:
|
||||
t0 = time.time()
|
||||
sock.connect((args.ip, args.port))
|
||||
log_success(f"[{thread_name}] 连接成功")
|
||||
local = sock.getsockname()
|
||||
remote = sock.getpeername()
|
||||
t_conn = (time.time() - t0) * 1000
|
||||
log_success(f"[{thread_name}] 连接成功 "
|
||||
f"(local={local[0]}:{local[1]}, "
|
||||
f"remote={remote[0]}:{remote[1]}, "
|
||||
f"{t_conn:.0f}ms)")
|
||||
|
||||
barrier.wait()
|
||||
|
||||
success_count = 0
|
||||
fail_count = 0
|
||||
|
||||
for i in range(args.count):
|
||||
try:
|
||||
# 发送
|
||||
test_data = f"T{client_id:02d}_{i:04d}".encode()
|
||||
t_send = time.time()
|
||||
sock.send(test_data)
|
||||
|
||||
# 接收
|
||||
data = sock.recv(100)
|
||||
if data and data == test_data:
|
||||
success_count += 1
|
||||
t_rtt = (time.time() - t_send) * 1000
|
||||
if data:
|
||||
if data == test_data:
|
||||
success_count += 1
|
||||
log_info(f"[{thread_name}] #{i}: OK ({t_rtt:.0f}ms)")
|
||||
else:
|
||||
fail_count += 1
|
||||
log_warn(f"[{thread_name}] #{i}: MISMATCH "
|
||||
f"want={len(test_data)}B/{test_data[:20].decode(errors='replace')!r} "
|
||||
f"got={len(data)}B/{data[:20].decode(errors='replace')!r}"
|
||||
f" ({t_rtt:.0f}ms)")
|
||||
else:
|
||||
fail_count += 1
|
||||
log_warn(f"[{thread_name}] #{i}: recv returned empty ({t_rtt:.0f}ms)")
|
||||
|
||||
except socket.timeout:
|
||||
fail_count += 1
|
||||
t_elapsed = (time.time() - t_send) * 1000
|
||||
log_error(f"[{thread_name}] #{i}: TIMEOUT ({t_elapsed:.0f}ms)")
|
||||
except ConnectionResetError:
|
||||
fail_count += 1
|
||||
t_elapsed = (time.time() - t_send) * 1000
|
||||
log_error(f"[{thread_name}] #{i}: RESET by peer ({t_elapsed:.0f}ms)")
|
||||
except Exception as e:
|
||||
fail_count += 1
|
||||
t_elapsed = (time.time() - t_send) * 1000
|
||||
log_error(f"[{thread_name}] #{i}: {type(e).__name__}: {e} ({t_elapsed:.0f}ms)")
|
||||
|
||||
results[client_id] = (success_count, fail_count)
|
||||
log_info(f"[{thread_name}] 完成: 成功 {success_count}, 失败 {fail_count}")
|
||||
t_total = (time.time() - t0) * 1000
|
||||
log_info(f"[{thread_name}] 完成: 成功 {success_count}, 失败 {fail_count} ({t_total:.0f}ms)")
|
||||
|
||||
except Exception as e:
|
||||
log_error(f"[{thread_name}] 错误: {e}")
|
||||
log_error(f"[{thread_name}] 连接失败: {type(e).__name__}: {e}")
|
||||
results[client_id] = (0, args.count)
|
||||
finally:
|
||||
sock.close()
|
||||
@@ -575,12 +604,13 @@ def stress_test(args):
|
||||
|
||||
results = {}
|
||||
threads = []
|
||||
barrier = threading.Barrier(num_clients)
|
||||
|
||||
# 创建并启动线程
|
||||
for i in range(num_clients):
|
||||
thread = threading.Thread(
|
||||
target=stress_test_thread,
|
||||
args=(i, args, results)
|
||||
args=(i, args, results, barrier)
|
||||
)
|
||||
threads.append(thread)
|
||||
thread.start()
|
||||
|
||||
Reference in New Issue
Block a user