Buffered output. Added net.monitor() and node.readline()
Moved monitor() and readline() into net.py and node.py respectively, which will hopefully be useful for monitoring large sets of hosts, as is done in udpbwtest.py. Changed iperf to use interactive command infrastructure (such as it is), which may make it more reliable. Hopefully it's a bit clearer as well, although it is slightly more complicated.
This commit is contained in:
+3
-34
@@ -18,7 +18,6 @@ various Mininet configurations, this example:
|
||||
|
||||
import os
|
||||
import re
|
||||
import select
|
||||
import sys
|
||||
from time import time
|
||||
|
||||
@@ -30,37 +29,6 @@ from mininet.node import KernelSwitch
|
||||
from mininet.topolib import TreeTopo
|
||||
from mininet.util import quietRun
|
||||
|
||||
# Some useful stuff: buffered readline and host monitoring
|
||||
|
||||
def readline( host, buf ):
|
||||
"Read a line from a host, buffering with buffer."
|
||||
buf += host.read( 1024 )
|
||||
if '\n' not in buf:
|
||||
return None, buf
|
||||
pos = buf.find( '\n' )
|
||||
line = buf[ 0 : pos ]
|
||||
rest = buf[ pos + 1: ]
|
||||
return line, rest
|
||||
|
||||
def monitor( hosts, seconds ):
|
||||
"Monitor a set of hosts and yield their output."
|
||||
poller = select.poll()
|
||||
Node = hosts[ 0 ] # so we can call class method fdToNode
|
||||
buffers = {}
|
||||
for host in hosts:
|
||||
poller.register( host.stdout )
|
||||
buffers[ host ] = ''
|
||||
quitTime = time() + seconds
|
||||
while time() < quitTime:
|
||||
ready = poller.poll()
|
||||
for fd, event in ready:
|
||||
host = Node.fdToNode( fd )
|
||||
if event & select.POLLIN:
|
||||
line, buffers[ host ] = readline( host, buffers[ host ] )
|
||||
if line:
|
||||
yield host, line
|
||||
yield None, ''
|
||||
|
||||
# bwtest support
|
||||
|
||||
def parsebwtest( line,
|
||||
@@ -111,8 +79,9 @@ def udpbwtest( net, seconds ):
|
||||
print
|
||||
results = {}
|
||||
print "*** Monitoring hosts"
|
||||
output = monitor( hosts, seconds )
|
||||
while True:
|
||||
output = net.monitor( hosts )
|
||||
quitTime = time() + seconds
|
||||
while time() < quitTime:
|
||||
host, line = output.next()
|
||||
if host is None:
|
||||
break
|
||||
|
||||
+33
-10
@@ -85,6 +85,7 @@ method may be called to shut down the network.
|
||||
|
||||
import os
|
||||
import re
|
||||
import select
|
||||
import signal
|
||||
from time import sleep
|
||||
|
||||
@@ -382,6 +383,25 @@ class Mininet( object ):
|
||||
self.stop()
|
||||
return result
|
||||
|
||||
def monitor( self, hosts=None ):
|
||||
"""Monitor a set of hosts (or all hosts by default),
|
||||
and return their output, a line at a time.
|
||||
returns: host, line"""
|
||||
if hosts is None:
|
||||
hosts = self.hosts
|
||||
poller = select.poll()
|
||||
Node = hosts[ 0 ] # so we can call class method fdToNode
|
||||
for host in hosts:
|
||||
poller.register( host.stdout )
|
||||
while True:
|
||||
ready = poller.poll()
|
||||
for fd, event in ready:
|
||||
host = Node.fdToNode( fd )
|
||||
if event & select.POLLIN:
|
||||
line = host.readline()
|
||||
if line:
|
||||
yield host, line
|
||||
|
||||
@staticmethod
|
||||
def _parsePing( pingOutput ):
|
||||
"Parse ping output and return packets sent, received."
|
||||
@@ -460,10 +480,10 @@ class Mininet( object ):
|
||||
hosts = [ self.hosts[ 0 ], self.hosts[ -1 ] ]
|
||||
else:
|
||||
assert len( hosts ) == 2
|
||||
host0, host1 = hosts
|
||||
client, server = hosts
|
||||
output( '*** Iperf: testing ' + l4Type + ' bandwidth between ' )
|
||||
output( "%s and %s\n" % ( host0.name, host1.name ) )
|
||||
host0.cmd( 'killall -9 iperf' )
|
||||
output( "%s and %s\n" % ( client.name, server.name ) )
|
||||
server.cmd( 'killall -9 iperf' )
|
||||
iperfArgs = 'iperf '
|
||||
bwArgs = ''
|
||||
if l4Type == 'UDP':
|
||||
@@ -471,14 +491,17 @@ class Mininet( object ):
|
||||
bwArgs = '-b ' + udpBw + ' '
|
||||
elif l4Type != 'TCP':
|
||||
raise Exception( 'Unexpected l4 type: %s' % l4Type )
|
||||
server = host0.cmd( iperfArgs + '-s &' )
|
||||
debug( '%s\n' % server )
|
||||
client = host1.cmd( iperfArgs + '-t 5 -c ' + host0.IP() + ' ' +
|
||||
server.sendCmd( iperfArgs + '-s', printPid=True )
|
||||
servout = ''
|
||||
while server.lastPid is None:
|
||||
servout += server.monitor()
|
||||
cliout = client.cmd( iperfArgs + '-t 5 -c ' + server.IP() + ' ' +
|
||||
bwArgs )
|
||||
debug( '%s\n' % client )
|
||||
server = host0.cmd( 'killall -9 iperf' )
|
||||
debug( '%s\n' % server )
|
||||
result = [ self._parseIperf( server ), self._parseIperf( client ) ]
|
||||
debug( 'Client output: %s\n' % cliout )
|
||||
server.sendInt()
|
||||
servout += server.waitOutput()
|
||||
debug( 'Server output: %s\n' % servout )
|
||||
result = [ self._parseIperf( servout ), self._parseIperf( cliout ) ]
|
||||
if l4Type == 'UDP':
|
||||
result.insert( 0, udpBw )
|
||||
output( '*** Results: %s\n' % result )
|
||||
|
||||
+30
-12
@@ -94,11 +94,7 @@ class Node( object ):
|
||||
self.defaultMAC = defaultMAC
|
||||
self.lastCmd = None
|
||||
self.lastPid = None
|
||||
# Grab PID
|
||||
self.waiting = True
|
||||
while self.lastPid is None:
|
||||
self.monitor()
|
||||
self.pid = self.lastPid
|
||||
self.readbuf = ''
|
||||
self.waiting = False
|
||||
|
||||
@classmethod
|
||||
@@ -115,9 +111,30 @@ class Node( object ):
|
||||
|
||||
# Subshell I/O, commands and control
|
||||
def read( self, bytes ):
|
||||
"""Read from a node.
|
||||
bytes: maximum number of bytes to read"""
|
||||
return os.read( self.stdout.fileno(), bytes )
|
||||
"""Buffered read from node, non-blocking.
|
||||
bytes: maximum number of bytes to return"""
|
||||
count = len( self.readbuf )
|
||||
if count < bytes:
|
||||
data = os.read( self.stdout.fileno(), bytes - count )
|
||||
self.readbuf += data
|
||||
if bytes >= len( self.readbuf ):
|
||||
result = self.readbuf
|
||||
self.readbuf = ''
|
||||
else:
|
||||
result = self.readbuf[ :bytes ]
|
||||
self.readbuf = self.readbuf[ bytes: ]
|
||||
return result
|
||||
|
||||
def readline( self ):
|
||||
"""Buffered readline from node, non-blocking.
|
||||
returns: line (minus newline) or None"""
|
||||
self.readbuf += self.read( 1024 )
|
||||
if '\n' not in self.readbuf:
|
||||
return None
|
||||
pos = self.readbuf.find( '\n' )
|
||||
line = self.readbuf[ 0 : pos ]
|
||||
self.readbuf = self.readbuf[ pos + 1: ]
|
||||
return line
|
||||
|
||||
def write( self, data ):
|
||||
"""Write data to node.
|
||||
@@ -135,7 +152,8 @@ class Node( object ):
|
||||
|
||||
def waitReadable( self ):
|
||||
"Wait until node's output is readable."
|
||||
self.pollOut.poll()
|
||||
if len( self.readbuf ) == 0:
|
||||
self.pollOut.poll()
|
||||
|
||||
def sendCmd( self, cmd, printPid=False ):
|
||||
"""Send a command, followed by a command to echo a sentinel,
|
||||
@@ -155,10 +173,10 @@ class Node( object ):
|
||||
self.lastPid = None
|
||||
self.waiting = True
|
||||
|
||||
def sendInt( self ):
|
||||
def sendInt( self, sig=signal.SIGINT ):
|
||||
"Interrupt running command."
|
||||
if self.lastPid:
|
||||
os.kill( self.lastPid, signal.SIGINT )
|
||||
os.kill( self.lastPid, sig )
|
||||
|
||||
def monitor( self ):
|
||||
"""Monitor and return the output of a command.
|
||||
@@ -315,7 +333,7 @@ class Node( object ):
|
||||
def setDefaultRoute( self, intf ):
|
||||
"""Set the default route to go through intf.
|
||||
intf: string, interface name"""
|
||||
self.cmd( 'ip route flush' )
|
||||
self.cmd( 'ip route flush root 0/0' )
|
||||
return self.cmd( 'route add default ' + intf )
|
||||
|
||||
def IP( self, intf=None ):
|
||||
|
||||
Reference in New Issue
Block a user