diff --git a/examples/udpbwtest.py b/examples/udpbwtest.py index 9832d18..2a899a3 100755 --- a/examples/udpbwtest.py +++ b/examples/udpbwtest.py @@ -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 diff --git a/mininet/net.py b/mininet/net.py index f732c36..8fdc07c 100755 --- a/mininet/net.py +++ b/mininet/net.py @@ -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 ) diff --git a/mininet/node.py b/mininet/node.py index 8037f2c..18b65d1 100644 --- a/mininet/node.py +++ b/mininet/node.py @@ -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 ):