From ec7b211c4176279d4d767bc13380caceb5ee6595 Mon Sep 17 00:00:00 2001 From: Bob Lantz Date: Tue, 16 Mar 2010 14:46:37 -0700 Subject: [PATCH] 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. --- examples/udpbwtest.py | 37 +++---------------------------------- mininet/net.py | 43 +++++++++++++++++++++++++++++++++---------- mininet/node.py | 42 ++++++++++++++++++++++++++++++------------ 3 files changed, 66 insertions(+), 56 deletions(-) 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 ):