diff --git a/bin/mn b/bin/mn index bdb5ea3..d69b6b1 100755 --- a/bin/mn +++ b/bin/mn @@ -41,7 +41,8 @@ from functools import partial # Experimental! cluster edition prototype from mininet.examples.cluster import ( MininetCluster, RemoteHost, RemoteOVSSwitch, RemoteLink, - SwitchBinPlacer, RandomPlacer ) + SwitchBinPlacer, RandomPlacer, + ClusterCleanup ) from mininet.examples.clustercli import ClusterCLI PLACEMENT = { 'block': SwitchBinPlacer, 'random': RandomPlacer } @@ -281,6 +282,11 @@ class MininetRunner( object ): def begin( self ): "Create and run mininet." + if self.options.cluster: + servers = self.options.cluster.split( ',' ) + for server in servers: + ClusterCleanup.add( server ) + if self.options.clean: cleanup() exit() @@ -334,7 +340,7 @@ class MininetRunner( object ): warn( '*** WARNING: Experimental cluster mode!\n' '*** Using RemoteHost, RemoteOVSSwitch, RemoteLink\n' ) host, switch, link = RemoteHost, RemoteOVSSwitch, RemoteLink - Net = partial( MininetCluster, servers=cluster.split( ',' ), + Net = partial( MininetCluster, servers=servers, placement=PLACEMENT[ self.options.placement ] ) mn = Net( topo=topo, diff --git a/examples/cluster.py b/examples/cluster.py index 7f3c09b..aa7ddfd 100755 --- a/examples/cluster.py +++ b/examples/cluster.py @@ -82,6 +82,7 @@ from mininet.topolib import TreeTopo from mininet.util import quietRun, errRun, retry from mininet.examples.clustercli import CLI from mininet.log import setLogLevel, debug, info, error +from mininet.clean import addCleanupCallback from signal import signal, SIGINT, SIG_IGN from subprocess import Popen, PIPE, STDOUT @@ -89,9 +90,51 @@ import os from random import randrange import sys import re - +from itertools import groupby +from operator import attrgetter from distutils.version import StrictVersion + +def findUser(): + "Try to return logged-in (usually non-root) user" + return ( + # If we're running sudo + os.environ.get( 'SUDO_USER', False ) or + # Logged-in user (if we have a tty) + ( quietRun( 'who am i' ).split() or [ False ] )[ 0 ] or + # Give up and return effective user + quietRun( 'whoami' ) ) + + +class ClusterCleanup( object ): + "Cleanup callback" + + inited = False + serveruser = {} + + @classmethod + def add( cls, server, user='' ): + "Add an entry to server: user dict" + if not cls.inited: + addCleanupCallback( cls.cleanup ) + if not user: + user = findUser() + cls.serveruser[ server ] = user + + @classmethod + def cleanup( cls ): + "Clean up" + info( '*** Cleaning up cluster\n' ) + for server, user in cls.serveruser.iteritems(): + if server == 'localhost': + # Handled by mininet.clean.cleanup() + continue + else: + cmd = [ 'su', user, '-c', + 'ssh %s@%s sudo mn -c' % ( user, server ) ] + info( cmd, '\n' ) + info( quietRun( cmd ) ) + # BL note: so little code is required for remote nodes, # we will probably just want to update the main Node() # class to enable it for remote access! However, there @@ -125,7 +168,8 @@ class RemoteMixin( object ): self.server = server if server else 'localhost' self.serverIP = ( serverIP if serverIP else self.findServerIP( self.server ) ) - self.user = user if user else self.findUser() + self.user = user if user else findUser() + ClusterCleanup.add( server=server, user=user ) if controlPath is True: # Set a default control path for shared SSH connections controlPath = '/tmp/mn-%r@%h:%p' @@ -148,17 +192,6 @@ class RemoteMixin( object ): self.shell, self.pid = None, None super( RemoteMixin, self ).__init__( name, **kwargs ) - @staticmethod - def findUser(): - "Try to return logged-in (usually non-root) user" - return ( - # If we're running sudo - os.environ.get( 'SUDO_USER', False ) or - # Logged-in user (if we have a tty) - ( quietRun( 'who am i' ).split() or [ False ] )[ 0 ] or - # Give up and return effective user - quietRun( 'whoami' ) ) - # Determine IP address of local host _ipMatchRegex = re.compile( r'\d+\.\d+\.\d+\.\d+' ) @@ -244,7 +277,7 @@ class RemoteMixin( object ): # Drop privileges cmd = [ 'sudo', '-E', '-u', self.user ] + cmd params.update( preexec_fn=self._ignoreSignal ) - debug( '_popen', ' '.join(cmd), params ) + debug( '_popen', cmd, '\n' ) popen = super( RemoteMixin, self )._popen( cmd, **params ) return popen @@ -257,13 +290,6 @@ class RemoteMixin( object ): kwargs.update( moveIntfFn=RemoteLink.moveIntf ) return super( RemoteMixin, self).addIntf( *args, **kwargs ) - def cleanup( self ): - "Help python collect its garbage." - # Intfs may end up in root NS - for intfName in self.intfNames(): - if self.name in intfName: - self.rcmd( 'ip link del ' + intfName ) - self.shell = None class RemoteNode( RemoteMixin, Node ): "A node on a remote server" @@ -280,6 +306,11 @@ class RemoteOVSSwitch( RemoteMixin, OVSSwitch ): OVSVersions = {} + def __init__( self, *args, **kwargs ): + # No batch startup yet + kwargs.update( batch=True ) + super( RemoteOVSSwitch, self ).__init__( *args, **kwargs ) + def isOldOVS( self ): "Is remote switch using an old OVS version?" cls = type( self ) @@ -293,9 +324,24 @@ class RemoteOVSSwitch( RemoteMixin, OVSSwitch ): StrictVersion( '1.10' ) ) @classmethod - def batchShutdown( cls, *_args, **_kwargs ): - "Not implemented yet" - return False + def batchStartup( cls, switches, **_kwargs ): + "Start up switches in per-server batches" + for server, switchGroup in groupby( switches, attrgetter( 'server' ) ): + info( '(%s)' % server ) + group = tuple( switchGroup ) + switch = group[ 0 ] + OVSSwitch.batchStartup( group, run=switch.cmd ) + return switches + + @classmethod + def batchShutdown( cls, switches, **_kwargs ): + "Stop switches in per-server batches" + for server, switchGroup in groupby( switches, attrgetter( 'server' ) ): + info( '(%s)' % server ) + group = tuple( switchGroup ) + switch = group[ 0 ] + OVSSwitch.batchShutdown( group, run=switch.rcmd ) + return switches class RemoteLink( Link ): @@ -315,6 +361,7 @@ class RemoteLink( Link ): def stop( self ): "Stop this link" + Link.stop( self ) if self.tunnel: self.tunnel.terminate() self.tunnel = None @@ -626,7 +673,7 @@ class MininetCluster( Mininet ): if not self.serverIP: self.serverIP = { server: RemoteMixin.findServerIP( server ) for server in self.servers } - self.user = params.pop( 'user', RemoteMixin.findUser() ) + self.user = params.pop( 'user', findUser() ) if params.pop( 'precheck' ): self.precheck() self.connections = {} diff --git a/examples/controlnet.py b/examples/controlnet.py index 9397188..50f989c 100755 --- a/examples/controlnet.py +++ b/examples/controlnet.py @@ -27,11 +27,17 @@ from mininet.log import setLogLevel, info class DataController( Controller ): """Data Network Controller. - patched to avoid checkListening error""" + patched to avoid checkListening error and to delete intfs""" + def checkListening( self ): "Ignore spurious error" pass + def stop( self, *args, **kwargs ): + "Make sure intfs are deleted" + kwargs.update( deleteIntfs=True ) + super( Controller, self ).stop( *args, **kwargs ) + class MininetFacade( object ): """Mininet object facade that allows a single CLI to talk to one or more networks""" diff --git a/examples/sshd.py b/examples/sshd.py index 39107ec..3fbd182 100755 --- a/examples/sshd.py +++ b/examples/sshd.py @@ -39,7 +39,7 @@ def connectToRootNS( network, switch, ip, routes ): routes: host networks to route to""" # Create a node in root namespace and link to switch 0 root = Node( 'root', inNamespace=False ) - intf = Link( root, switch ).intf1 + intf = network.addLink( root, switch ).intf1 root.setIP( ip, intf=intf ) # Start network that now includes link to root namespace network.start() diff --git a/examples/tree1024.py b/examples/tree1024.py index 9397131..d5c25c2 100755 --- a/examples/tree1024.py +++ b/examples/tree1024.py @@ -9,10 +9,10 @@ and running sysctl -p. Check util/sysctl_addon. from mininet.cli import CLI from mininet.log import setLogLevel -from mininet.node import OVSKernelSwitch +from mininet.node import OVSSwitch from mininet.topolib import TreeNet if __name__ == '__main__': setLogLevel( 'info' ) - network = TreeNet( depth=2, fanout=32, switch=OVSKernelSwitch ) + network = TreeNet( depth=2, fanout=32, switch=OVSSwitch ) network.run( CLI, network ) diff --git a/mininet/clean.py b/mininet/clean.py index 4761721..707452e 100755 --- a/mininet/clean.py +++ b/mininet/clean.py @@ -38,60 +38,80 @@ def killprocs( pattern ): else: break -def cleanup(): - """Clean up junk which might be left over from old runs; - do fast stuff before slow dp and link removal!""" +class Cleanup( object ): + "Wrapper for cleanup()" - info("*** Removing excess controllers/ofprotocols/ofdatapaths/pings/noxes" - "\n") - zombies = 'controller ofprotocol ofdatapath ping nox_core lt-nox_core ' - zombies += 'ovs-openflowd ovs-controller udpbwtest mnexec ivs' - # Note: real zombie processes can't actually be killed, since they - # are already (un)dead. Then again, - # you can't connect to them either, so they're mostly harmless. - # Send SIGTERM first to give processes a chance to shutdown cleanly. - sh( 'killall ' + zombies + ' 2> /dev/null' ) - time.sleep( 1 ) - sh( 'killall -9 ' + zombies + ' 2> /dev/null' ) + callbacks = [] - # And kill off sudo mnexec - sh( 'pkill -9 -f "sudo mnexec"') + @classmethod + def cleanup( cls): + """Clean up junk which might be left over from old runs; + do fast stuff before slow dp and link removal!""" - info( "*** Removing junk from /tmp\n" ) - sh( 'rm -f /tmp/vconn* /tmp/vlogs* /tmp/*.out /tmp/*.log' ) + info("*** Removing excess controllers/ofprotocols/ofdatapaths/pings/noxes" + "\n") + zombies = 'controller ofprotocol ofdatapath ping nox_core lt-nox_core ' + zombies += 'ovs-openflowd ovs-controller udpbwtest mnexec ivs' + # Note: real zombie processes can't actually be killed, since they + # are already (un)dead. Then again, + # you can't connect to them either, so they're mostly harmless. + # Send SIGTERM first to give processes a chance to shutdown cleanly. + sh( 'killall ' + zombies + ' 2> /dev/null' ) + time.sleep( 1 ) + sh( 'killall -9 ' + zombies + ' 2> /dev/null' ) - info( "*** Removing old X11 tunnels\n" ) - cleanUpScreens() + # And kill off sudo mnexec + sh( 'pkill -9 -f "sudo mnexec"') - info( "*** Removing excess kernel datapaths\n" ) - dps = sh( "ps ax | egrep -o 'dp[0-9]+' | sed 's/dp/nl:/'" ).splitlines() - for dp in dps: - if dp: - sh( 'dpctl deldp ' + dp ) + info( "*** Removing junk from /tmp\n" ) + sh( 'rm -f /tmp/vconn* /tmp/vlogs* /tmp/*.out /tmp/*.log' ) - info( "*** Removing OVS datapaths" ) - dps = sh("ovs-vsctl --timeout=1 list-br").strip().splitlines() - if dps: - sh( "ovs-vsctl " + " -- ".join( "--if-exists del-br " + dp - for dp in dps if dp ) ) - # And in case the above didn't work... - dps = sh("ovs-vsctl --timeout=1 list-br").strip().splitlines() - for dp in dps: - sh( 'ovs-vsctl del-br ' + dp ) + info( "*** Removing old X11 tunnels\n" ) + cleanUpScreens() - info( "*** Removing all links of the pattern foo-ethX\n" ) - links = sh( "ip link show | " - "egrep -o '([-_.[:alnum:]]+-eth[[:digit:]]+)'" ).splitlines() - for link in links: - if link: - sh( "ip link del " + link ) + info( "*** Removing excess kernel datapaths\n" ) + dps = sh( "ps ax | egrep -o 'dp[0-9]+' | sed 's/dp/nl:/'" ).splitlines() + for dp in dps: + if dp: + sh( 'dpctl deldp ' + dp ) - info( "*** Killing stale mininet node processes\n" ) - killprocs( 'mininet:' ) + info( "*** Removing OVS datapaths" ) + dps = sh("ovs-vsctl --timeout=1 list-br").strip().splitlines() + if dps: + sh( "ovs-vsctl " + " -- ".join( "--if-exists del-br " + dp + for dp in dps if dp ) ) + # And in case the above didn't work... + dps = sh("ovs-vsctl --timeout=1 list-br").strip().splitlines() + for dp in dps: + sh( 'ovs-vsctl del-br ' + dp ) - info( "*** Shutting down stale tunnels\n" ) - killprocs( 'Tunnel=Ethernet' ) - killprocs( '.ssh/mn') - sh( 'rm -f ~/.ssh/mn/*' ) + info( "*** Removing all links of the pattern foo-ethX\n" ) + links = sh( "ip link show | " + "egrep -o '([-_.[:alnum:]]+-eth[[:digit:]]+)'" ).splitlines() + for link in links: + if link: + sh( "ip link del " + link ) - info( "*** Cleanup complete.\n" ) + info( "*** Killing stale mininet node processes\n" ) + killprocs( 'mininet:' ) + + info( "*** Shutting down stale tunnels\n" ) + killprocs( 'Tunnel=Ethernet' ) + killprocs( '.ssh/mn') + sh( 'rm -f ~/.ssh/mn/*' ) + + # Call any additional cleanup code if necessary + for callback in cls.callbacks: + callback() + + info( "*** Cleanup complete.\n" ) + + @classmethod + def addCleanupCallback( cls, callback ): + "Add cleanup callback" + if callback not in cls.callbacks: + cls.callbacks.append( callback ) + + +cleanup = Cleanup.cleanup +addCleanupCallback = Cleanup.addCleanupCallback diff --git a/mininet/link.py b/mininet/link.py index af895a1..4a1a683 100644 --- a/mininet/link.py +++ b/mininet/link.py @@ -197,9 +197,10 @@ class Intf( object ): def delete( self ): "Delete interface" self.cmd( 'ip link del ' + self.name ) - if self.node.inNamespace: - # Link may have been dumped into root NS - quietRun( 'ip link del ' + self.name ) + # We used to do this, but it slows us down: + # if self.node.inNamespace: + # Link may have been dumped into root NS + # quietRun( 'ip link del ' + self.name ) def status( self ): "Return intf status as a string" diff --git a/mininet/net.py b/mininet/net.py index bb7b66d..7a01ecb 100755 --- a/mininet/net.py +++ b/mininet/net.py @@ -414,7 +414,12 @@ class Mininet( object ): info( '\n*** Adding switches:\n' ) for switchName in topo.switches(): - self.addSwitch( switchName, **topo.nodeInfo( switchName) ) + # A bit ugly: add batch parameter if appropriate + params = topo.nodeInfo( switchName) + cls = params.get( 'cls', self.switch ) + if hasattr( cls, 'batchStartup' ): + params.setdefault( 'batch', True ) + self.addSwitch( switchName, **params ) info( switchName + ' ' ) info( '\n*** Adding links:\n' ) @@ -481,6 +486,13 @@ class Mininet( object ): for switch in self.switches: info( switch.name + ' ') switch.start( self.controllers ) + started = {} + for swclass, switches in groupby( + sorted( self.switches, key=type ), type ): + switches = tuple( switches ) + if hasattr( swclass, 'batchStartup' ): + success = swclass.batchStartup( switches ) + started.update( { s: s for s in success } ) info( '\n' ) if self.waitConn: self.waitConnected() @@ -505,9 +517,9 @@ class Mininet( object ): for swclass, switches in groupby( sorted( self.switches, key=type ), type ): switches = tuple( switches ) - if ( hasattr( swclass, 'batchShutdown' ) and - swclass.batchShutdown( switches ) ): - stopped.update( { s: s for s in switches } ) + if hasattr( swclass, 'batchShutdown' ): + success = swclass.batchShutdown( switches ) + stopped.update( { s: s for s in success } ) for switch in self.switches: info( switch.name + ' ' ) if switch not in stopped: diff --git a/mininet/node.py b/mininet/node.py index 1e0e121..c602919 100644 --- a/mininet/node.py +++ b/mininet/node.py @@ -157,8 +157,7 @@ class Node( object ): break self.pollOut.poll() self.waiting = False - self.cmd( 'stty -echo' ) - self.cmd( 'set +m' ) + self.cmd( 'stty -echo; set +m' ) def mountPrivateDirs( self ): "mount private directories" @@ -194,10 +193,11 @@ class Node( object ): def cleanup( self ): "Help python collect its garbage." + # We used to do this, but it slows us down: # Intfs may end up in root NS - for intfName in self.intfNames(): - if self.name in intfName: - quietRun( 'ip link del ' + intfName ) + # for intfName in self.intfNames(): + # if self.name in intfName: + # quietRun( 'ip link del ' + intfName ) self.shell = None # Subshell I/O, commands and control @@ -258,9 +258,9 @@ class Node( object ): """Send a command, followed by a command to echo a sentinel, and return without waiting for the command to complete. args: command and arguments, or string - printPid: print command's PID?""" + printPid: print command's PID? (False)""" assert self.shell and not self.waiting - printPid = kwargs.get( 'printPid', True ) + printPid = kwargs.get( 'printPid', False ) # Allow sendCmd( [ list ] ) if len( args ) == 1 and isinstance( args[ 0 ], list ): cmd = args[ 0 ] @@ -1058,7 +1058,7 @@ class OVSSwitch( Switch ): def __init__( self, name, failMode='secure', datapath='kernel', inband=False, protocols=None, - reconnectms=1000, stp=False, **params ): + reconnectms=1000, stp=False, batch=False, **params ): """name: name for switch failMode: controller loss behavior (secure|open) datapath: userspace or kernel mode (kernel|user) @@ -1066,7 +1066,8 @@ class OVSSwitch( Switch ): protocols: use specific OpenFlow version(s) (e.g. OpenFlow13) Unspecified (or old OVS version) uses OVS default reconnectms: max reconnect timeout in ms (0/None for default) - stp: enable STP (False, requires failMode=standalone)""" + stp: enable STP (False, requires failMode=standalone) + batch: enable batch startup (False)""" Switch.__init__( self, name, **params ) self.failMode = failMode self.datapath = datapath @@ -1075,6 +1076,8 @@ class OVSSwitch( Switch ): self.reconnectms = reconnectms self.stp = stp self._uuids = [] # controller UUIDs + self.batch = batch + self.commands = [] # saved commands for batch startup @classmethod def setup( cls ): @@ -1104,18 +1107,18 @@ class OVSSwitch( Switch ): return ( StrictVersion( cls.OVSVersion ) < StrictVersion( '1.10' ) ) - @classmethod - def batchShutdown( cls, switches ): - "Call ovs-vsctl del-br on all OVSSwitches in a list" - quietRun( 'ovs-vsctl ' + - ' -- '.join( '--if-exists del-br %s' % s - for s in switches ) ) - return True - def dpctl( self, *args ): "Run ovs-ofctl command" return self.cmd( 'ovs-ofctl', args[ 0 ], self, *args[ 1: ] ) + def vsctl( self, *args, **kwargs ): + "Run ovs-vsctl command (or queue for later execution)" + if self.batch: + cmd = ' '.join( str( arg ).strip() for arg in args ) + self.commands.append( cmd ) + else: + return self.cmd( 'ovs-vsctl', *args, **kwargs ) + @staticmethod def TCReapply( intf ): """Unfortunately OVS and Mininet are fighting @@ -1126,13 +1129,13 @@ class OVSSwitch( Switch ): def attach( self, intf ): "Connect a data port" - self.cmd( 'ovs-vsctl add-port', self, intf ) + self.vsctl( 'add-port', self, intf ) self.cmd( 'ifconfig', intf, 'up' ) self.TCReapply( intf ) def detach( self, intf ): "Disconnect a data port" - self.cmd( 'ovs-vsctl del-port', self, intf ) + self.vsctl( 'del-port', self, intf ) def controllerUUIDs( self, update=False ): """Return ovsdb UUIDs for our controllers @@ -1150,84 +1153,109 @@ class OVSSwitch( Switch ): def connected( self ): "Are we connected to at least one of our controllers?" for uuid in self.controllerUUIDs(): - if 'true' in self.cmd( 'ovs-vsctl -- get Controller', - uuid, 'is_connected' ): + if 'true' in self.vsctl( '-- get Controller', + uuid, 'is_connected' ): return True return self.failMode == 'standalone' - @staticmethod - def patchOpts( intf ): - "Return OVS patch port options (if any) for intf" - if not isinstance( intf, OVSIntf ): - # Ignore if it's not a patch link - return '' - intf1, intf2 = intf.link.intf1, intf.link.intf2 - peer = intf1 if intf1 != intf else intf2 - return ( '-- set Interface %s type=patch ' - '-- set Interface %s options:peer=%s ' % - ( intf, intf, peer ) ) + def intfOpts( self, intf ): + "Return OVS interface options for intf" + opts = '' + if not self.isOldOVS(): + # ofport_request is not supported on old OVS + opts += ' ofport_request=%s' % self.ports[ intf ] + # Patch ports don't work well with old OVS + if isinstance( intf, OVSIntf ): + intf1, intf2 = intf.link.intf1, intf.link.intf2 + peer = intf1 if intf1 != intf else intf2 + opts += ' type=patch options:peer=%s' % peer + return '' if not opts else ' -- set Interface %s' % intf + opts + + def bridgeOpts( self ): + "Return OVS bridge options" + opts = ( ' other_config:datapath-id=%s' % self.dpid + + ' fail_mode=%s' % self.failMode ) + if not self.inband: + opts += ' other-config:disable-in-band=true' + if self.datapath == 'user': + opts += ' datapath_type=netdev' + if self.protocols and not self.isOldOVS(): + opts += ' protocols=%s' % ( self, self.protocols ) + if self.stp and self.failMode == 'standalone': + opts += ' stp_enable=true' % self + return opts - # pylint: disable=too-many-branches def start( self, controllers ): "Start up a new OVS OpenFlow switch using ovs-vsctl" if self.inNamespace: raise Exception( 'OVS kernel switch does not work in a namespace' ) int( self.dpid, 16 ) # DPID must be a hex string - # Interfaces and controllers - intfs = ' '.join( '-- add-port %s %s ' % ( self, intf ) + - '-- set Interface %s ' % intf + - 'ofport_request=%s ' % self.ports[ intf ] - + self.patchOpts( intf ) - for intf in self.intfList() - if self.ports[ intf ] and not intf.IP() ) - clist = ' '.join( '%s:%s:%d' % ( c.protocol, c.IP(), c.port ) - for c in controllers ) + # Command to add interfaces + intfs = ''.join( ' -- add-port %s %s' % ( self, intf ) + + self.intfOpts( intf ) + for intf in self.intfList() + if self.ports[ intf ] and not intf.IP() ) + # Command to create controller entries + clist = [ ( self.name + c.name, '%s:%s:%d' % + ( c.protocol, c.IP(), c.port ) ) + for c in controllers ] if self.listenPort: - clist += ' ptcp:%s' % self.listenPort - # Construct big ovs-vsctl command for new versions of OVS - if not self.isOldOVS(): - cmd = ( 'ovs-vsctl --if-exists del-br %s ' % self + - '-- add-br %s ' % self + - '-- set Bridge %s ' % self + - 'other_config:datapath-id=%s ' % self.dpid + - '-- set-fail-mode %s %s ' % ( self, self.failMode ) + - intfs + - '-- set-controller %s %s ' % ( self, clist ) ) - # Construct ovs-vsctl commands for old versions of OVS - else: - # Annoyingly, --if-exists option seems not to work - self.cmd( 'ovs-vsctl del-br', self ) - self.cmd( 'ovs-vsctl add-br', self ) - for intf in self.intfList(): - if not intf.IP(): - self.cmd( 'ovs-vsctl add-port', self, intf ) - cmd = ( 'ovs-vsctl set Bridge %s ' % self + - 'other_config:datapath-id=%s ' % self.dpid + - '-- set-fail-mode %s %s ' % ( self, self.failMode ) + - '-- set-controller %s %s ' % ( self, clist ) ) - if not self.inband: - cmd += ( '-- set bridge %s ' - 'other-config:disable-in-band=true ' % self ) - if self.datapath == 'user': - cmd += '-- set bridge %s datapath_type=netdev ' % self - if self.protocols and not self.isOldOVS(): - cmd += '-- set bridge %s protocols=%s ' % ( self, self.protocols ) - if self.stp and self.failMode == 'standalone': - cmd += '-- set bridge %s stp_enable=true ' % self - # Do it!! - self.cmd( cmd ) - # Reconnect quickly to controllers (1s vs. 15s max_backoff) + clist.append( ( self.name + '-listen', + 'ptcp:%s' % self.listenPort ) ) + ccmd = '-- --id=@%s create Controller target=\\"%s\\"' if self.reconnectms: - uuids = [ '-- set Controller %s max_backoff=%d' % - ( uuid, self.reconnectms ) - for uuid in self.controllerUUIDs() ] - if uuids: - self.cmd( 'ovs-vsctl', *uuids ) + ccmd += ' max_backoff=%d' % self.reconnectms + cargs = ' '.join( ccmd % ( name, target ) + for name, target in clist ) + # Controller ID list + cids = ','.join( '@%s' % name for name, _target in clist ) + # Try to delete any existing bridges with the same name + if not self.isOldOVS(): + cargs += ' -- --if-exists del-br %s' % self + # One ovs-vsctl command to rule them all! + self.vsctl( cargs + + ' -- add-br %s' % self + + ' -- set bridge %s controller=[%s]' % ( self, cids ) + + self.bridgeOpts() + + intfs ) # If necessary, restore TC config overwritten by OVS - for intf in self.intfList(): - self.TCReapply( intf ) - # pylint: enable=too-many-branches + if not self.batch: + for intf in self.intfList(): + self.TCReapply( intf ) + + # This should be ~ int( quietRun( 'getconf ARG_MAX' ) ), + # but the real limit seems to be much lower + argmax = 128000 + + @classmethod + def batchStartup( cls, switches, run=errRun ): + """Batch startup for OVS + switches: switches to start up + run: function to run commands (errRun)""" + info( '...' ) + cmds = 'ovs-vsctl' + for switch in switches: + if switch.isOldOVS(): + # Ideally we'd optimize this also + run( 'ovs-vsctl del-br %s' % switch ) + for cmd in switch.commands: + cmd = cmd.strip() + # Don't exceed ARG_MAX + if len( cmds ) + len( cmd ) >= cls.argmax: + run( cmds, shell=True ) + cmds = 'ovs-vsctl' + cmds += ' ' + cmd + switch.cmds = [] + switch.batch = False + if cmds: + run( cmds, shell=True ) + # Reapply link config if necessary... + for switch in switches: + for intf in switch.intfs.itervalues(): + if isinstance( intf, TCIntf ): + intf.config( **intf.params ) + return switches def stop( self, deleteIntfs=True ): """Terminate OVS switch. @@ -1237,6 +1265,22 @@ class OVSSwitch( Switch ): self.cmd( 'ip link del', self ) super( OVSSwitch, self ).stop( deleteIntfs ) + @classmethod + def batchShutdown( cls, switches, run=errRun ): + "Shut down a list of OVS switches" + delcmd = 'del-br %s' + if switches and not switches[ 0 ].isOldOVS(): + delcmd = '--if-exists ' + delcmd + # First, delete them all from ovsdb + run( 'ovs-vsctl ' + + ' -- '.join( delcmd % s for s in switches ) ) + # Next, shut down all of the processes + pids = ' '.join( str( switch.pid ) for switch in switches ) + run( 'kill -HUP ' + pids ) + for switch in switches: + switch.shell = None + return switches + OVSKernelSwitch = OVSSwitch @@ -1285,6 +1329,7 @@ class IVSSwitch( Switch ): "Kill each IVS switch, to be waited on later in stop()" for switch in switches: switch.cmd( 'kill %ivs' ) + return switches def start( self, controllers ): "Start up a new IVS switch" @@ -1379,6 +1424,7 @@ class Controller( Node ): "Stop controller." self.cmd( 'kill %' + self.command ) self.cmd( 'wait %' + self.command ) + kwargs.update( deleteIntfs=False ) super( Controller, self ).stop( *args, **kwargs ) def IP( self, intf=None ): diff --git a/mininet/util.py b/mininet/util.py index a655f7d..1f5f64d 100644 --- a/mininet/util.py +++ b/mininet/util.py @@ -77,6 +77,7 @@ def errRun( *cmd, **kwargs ): cmd = [ str( arg ) for arg in cmd ] elif isinstance( cmd, list ) and shell: cmd = " ".join( arg for arg in cmd ) + debug( '*** errRun:', cmd, '\n' ) popen = Popen( cmd, stdout=PIPE, stderr=stderr, shell=shell ) # We use poll() because select() doesn't work with large fd numbers, # and thus communicate() doesn't work either @@ -113,6 +114,7 @@ def errRun( *cmd, **kwargs ): poller.unregister( fd ) returncode = popen.wait() + debug( out, err, returncode ) return out, err, returncode def errFail( *cmd, **kwargs ): @@ -188,7 +190,8 @@ def makeIntfPair( intf1, intf2, addr1=None, addr2=None, node1=None, node2=None, if cmdOutput == '': return True else: - error( "Error creating interface pair: %s " % cmdOutput ) + raise Exception( "Error creating interface pair (%s,%s): %s " % + ( intf1, intf2, cmdOutput ) ) return False def retry( retries, delaySecs, fn, *args, **keywords ): @@ -533,19 +536,28 @@ def customConstructor( constructors, argStr ): raise Exception( "error: %s is unknown - please specify one of %s" % ( cname, constructors.keys() ) ) - def customized( name, *args, **params ): - "Customized constructor, useful for Node, Link, and other classes" - params = params.copy() - params.update( kwargs ) - if not newargs: - return constructor( name, *args, **params ) - if args: - warn( 'warning: %s replacing %s with %s\n' % ( - constructor, args, newargs ) ) - return constructor( name, *newargs, **params ) + if not newargs and not kwargs: + return constructor - customized.__name__ = 'customConstructor(%s)' % argStr - return customized + if not isinstance( constructor, type ): + raise Exception( "error: invalid arguments %s" % argStr ) + + # Return a customized subclass + cls = constructor + class CustomClass( cls ): + "Customized subclass, useful for Node, Link, and other classes" + def __init__( self, name, *args, **params ): + params = params.copy() + params.update( kwargs ) + if not newargs: + return cls.__init__( self, name, *args, **params ) + if args: + warn( 'warning: %s replacing %s with %s\n' % + ( constructor, args, newargs ) ) + return cls.__init__( self, name, *newargs, **params ) + + CustomClass.__name__ = '%s%s' % ( cls.__name__, kwargs ) + return CustomClass def buildTopo( topos, topoStr ): """Create topology from string with format (object, arg1, arg2,...).