Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
129 changes: 84 additions & 45 deletions python-ecosys/iperf3/iperf3.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand All @@ -251,26 +259,40 @@ 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
udp_last_send = None
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", 987654321), addr)
s = socket.socket(ai[0], socket.SOCK_DGRAM)
s.bind(ai[-1])
data, addr = s.recvfrom(4)
s.sendto(struct.pack("<I", 987654321), addr)
s_data = [s]
udp_interval = 1000000 * 8 * param["len"] // param["bandwidth"]
udp = True
udp_packet_id = 0
Expand All @@ -287,10 +309,9 @@ def server_once():
# Read data, and wait for client to send TEST_END
poll = select.poll()
poll.register(s_ctrl, select.POLLIN)
if reverse:
poll.register(s_data, select.POLLOUT)
else:
poll.register(s_data, select.POLLIN)
for s in s_data:
poll.register(s, select.POLLOUT if reverse else select.POLLIN)
stream_bytes = [0] * len(s_data)
stats = Stats(param)
stats.start()
running = True
Expand All @@ -303,35 +324,47 @@ def server_once():
print(cmd_string.get(cmd, "UNKNOWN_COMMAND"))
if cmd == TEST_END:
running = False
elif pollable_is_sock(pollable, s_data):
udp_last_send, udp_packet_id = _transfer(
udp,
not reverse,
addr,
s_data,
data_buf,
udp_last_send,
udp_packet_id,
udp_interval,
stats,
)
stats.update()

# Need to continue writing so other side doesn't get blocked waiting for data
if reverse and not udp:
while True:
for pollable in poll.poll(0):
if pollable_is_sock(pollable, s_data):
s_data.send(data_buf)
continue
for i, s in enumerate(s_data):
if pollable_is_sock(pollable, s):
nb = stats.nb0
udp_last_send, udp_packet_id = _transfer(
udp,
not reverse,
addr,
s,
data_buf,
udp_last_send,
udp_packet_id,
udp_interval,
stats,
)
stream_bytes[i] += stats.nb0 - nb
break
else:
break
stats.update()

stats.stop()

# Ask to exchange results
s_ctrl.sendall(bytes([EXCHANGE_RESULTS]))

# Need to continue writing until the other side responds, because it may be blocked
# waiting for the rest of a block of data and not see the request until it gets it
if reverse and not udp:
waiting = True
while waiting:
for pollable in poll.poll():
if pollable_is_sock(pollable, s_ctrl):
waiting = False
continue
for s in s_data:
if pollable_is_sock(pollable, s):
try:
s.send(data_buf)
except OSError:
# The other side has closed this stream
poll.unregister(s)

# Get client results
n = struct.unpack(">I", recvn(s_ctrl, 4))[0]
results = recvn(s_ctrl, n)
Expand All @@ -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)
Expand All @@ -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()

Expand Down Expand Up @@ -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
Expand Down
2 changes: 1 addition & 1 deletion python-ecosys/iperf3/manifest.py
Original file line number Diff line number Diff line change
@@ -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")
Loading