diff --git a/python-ecosys/iperf3/iperf3.py b/python-ecosys/iperf3/iperf3.py index 05d69f774..5c8b1b176 100644 --- a/python-ecosys/iperf3/iperf3.py +++ b/python-ecosys/iperf3/iperf3.py @@ -71,7 +71,8 @@ def fmt_size(val, div): class Stats: def __init__(self, param): - self.pacing_timer_us = param["pacing_timer"] * 1000 + # iperf3 clients before 3.2 do not send the pacing timer. + self.pacing_timer_us = param.get("pacing_timer", 1000) * 1000 self.udp = param.get("udp", False) self.reverse = param.get("reverse", False) self.running = False @@ -217,23 +218,30 @@ def _transfer(udp, reverse, addr, s_data, buf, udp_last_send, udp_packet_id, udp stats.add_bytes(n) else: if reverse: - recvninto(s_data, buf) - n = len(buf) + # The socket is non-blocking, so this reads only what is available. The + # sender is not required to end the stream on a block boundary, so waiting + # here for a full block could block forever and miss TEST_END. + n = recvinto(s_data, buf) or 0 else: n = s_data.send(buf) stats.add_bytes(n) return udp_last_send, udp_packet_id +def _listen(ai, backlog): + s_listen = socket.socket(ai[0], socket.SOCK_STREAM) + s_listen.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1) + s_listen.bind(ai[-1]) + s_listen.listen(backlog) + return s_listen + + def server_once(): # Listen for a connection ai = socket.getaddrinfo("0.0.0.0", 5201) ai = ai[0] print("Server listening on", ai[-1]) - s_listen = socket.socket(ai[0], socket.SOCK_STREAM) - s_listen.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1) - s_listen.bind(ai[-1]) - s_listen.listen(1) + s_listen = _listen(ai, 1) s_ctrl, addr = s_listen.accept() # Read client's cookie @@ -251,15 +259,28 @@ def server_once(): if DEBUG: print(param) reverse = param.get("reverse", False) + parallel = param.get("parallel", 1) + + if parallel > 1: + # The client opens all of its streams at once, so listen again with room to + # queue them all. + s_listen.close() + s_listen = _listen(ai, parallel) # Ask to create streams s_ctrl.sendall(bytes([CREATE_STREAMS])) if param.get("tcp", False): - # Accept stream - s_data, addr = s_listen.accept() - print("Accepted connection:", addr) - recvn(s_data, COOKIE_SIZE) + # Accept streams, the client opens one connection per parallel stream + s_data = [] + for _ in range(parallel): + s, addr = s_listen.accept() + print("Accepted connection:", addr) + recvn(s, COOKIE_SIZE) + if not reverse: + # Receive without blocking, see _transfer(). + s.setblocking(False) + s_data.append(s) udp = False udp_packet_id = 0 udp_interval = None @@ -267,10 +288,11 @@ def server_once(): elif param.get("udp", False): # Close TCP connection and open UDP "connection" s_listen.close() - s_data = socket.socket(ai[0], socket.SOCK_DGRAM) - s_data.bind(ai[-1]) - data, addr = s_data.recvfrom(4) - s_data.sendto(struct.pack("I", recvn(s_ctrl, 4))[0] results = recvn(s_ctrl, n) @@ -348,15 +381,17 @@ def server_once(): "congestion_used": "cubic", "streams": [ { - "id": 1, - "bytes": stats.nb0, + # The reference implementation numbers its streams 1, 3, 4, 5, ... + "id": i + 2 if i else 1, + "bytes": stream_bytes[i], "retransmits": 0, "jitter": 0, "errors": 0, - "packets": stats.np0, + "packets": 0 if i else stats.np0, "start_time": 0, "end_time": ticks_diff(stats.t3, stats.t0) * 1e-6, } + for i in range(len(s_data)) ], } results = json.dumps(results) @@ -371,7 +406,8 @@ def server_once(): assert cmd == IPERF_DONE # Close all sockets - s_data.close() + for s in s_data: + s.close() s_ctrl.close() s_listen.close() @@ -488,6 +524,9 @@ def client(host, udp=False, reverse=False, bandwidth=10 * 1024 * 1024): s_data = socket.socket(ai[0], socket.SOCK_STREAM) s_data.connect(ai[-1]) s_data.sendall(cookie) + if reverse: + # Receive without blocking, see _transfer(). + s_data.setblocking(False) buf = bytearray(urandom(param["len"])) elif cmd == EXCHANGE_RESULTS: # Close data socket now that server knows we are finished, to prevent it flooding us diff --git a/python-ecosys/iperf3/manifest.py b/python-ecosys/iperf3/manifest.py index d63c2d6d1..8c1570d14 100644 --- a/python-ecosys/iperf3/manifest.py +++ b/python-ecosys/iperf3/manifest.py @@ -1,3 +1,3 @@ -metadata(version="0.1.5", pypi="iperf3", pypi_publish="uiperf3") +metadata(version="0.1.6", pypi="iperf3", pypi_publish="uiperf3") module("iperf3.py")