from teuthology import misc as teuthology
from cStringIO import StringIO
from teuthology.orchestra import run
-from teuthology.orchestra import remote as orchestra_remote
from teuthology import contextutil
from paramiko import SSHException
import socket
log = logging.getLogger(__name__)
+
def set_priority(interface):
# create a priority queueing discipline
return ['sudo', 'tc', 'qdisc', 'add', 'dev', interface, 'root', 'handle', '1:', 'prio']
+
def show_tc(interface):
# shows tc device present
return ['sudo', 'tc', 'qdisc', 'show', 'dev', interface]
+
def del_tc(interface):
return ['sudo', 'tc', 'qdisc', 'del', 'dev', interface, 'root']
+
def cmd_prefix(interface):
# prepare command to set delay
return cmd1, cmd2, cmd3
+
def static_delay(remote, host, interface, delay):
""" Sets a constant delay between two hosts to emulate network delays using tc qdisc and netem"""
log.info('Delay set on %s' % remote)
set_ip.extend(['%s' % ip, 'flowid', '2:1'])
remote.run(args=set_ip)
- remote.run(args=show_tc(interface), stdout=StringIO())
else:
# if the device is already created, only change the delay
log.info('Setting delay to %s' % delay)
change_delay.extend(['%s' % delay, '5ms', 'distribution', 'normal'])
remote.run(args=change_delay)
- remote.run(args=show_tc(interface), stdout=StringIO())
+
def variable_delay(remote, host, interface, delay_range=[]):
log.info('Delay set on %s' % remote)
set_ip.extend(['%s' % ip, 'flowid', '2:1'])
remote.run(args=set_ip)
- remote.run(args=show_tc(interface), stdout=StringIO())
else:
# if the device is already created, only change the delay
log.info('Setting varying delay')
change_delay.extend(['%s' % delay1, '%s' % delay2])
remote.run(args=change_delay)
- remote.run(args=show_tc(interface), stdout=StringIO())
def delete_dev(remote, interface):
stop_event = gevent.event.Event()
- def __init__(self, remote, host, interface, interval):
+ def __init__(self, ctx, remote, host, interface, interval):
+ self.ctx = ctx
self.remote = remote
self.host = host
self.interval = interval
set_ip.extend(['%s' % self.ip, 'action', 'drop'])
self.remote.run(args=set_ip)
-
def link_toggle(self):
- # For toggling packet drop and recovery in regular interval.
- # If interval is 5s, link is up for 5s and link is down for 5s
-
- while not self.stop_event.is_set():
- self.stop_event.wait(timeout=self.interval)
- # simulate link down
- try:
- self.packet_drop()
- log.info('link down')
- except SSHException as e:
- log.debug('Failed to run command')
-
- self.stop_event.wait(timeout=self.interval)
- # if qdisc exist,delete it.
- try:
- delete_dev(self.remote, self.interface)
- log.info('link up')
- except SSHException as e:
- log.debug('Failed to run command')
-
- def begin(self):
+
+ """
+ For toggling packet drop and recovery in regular interval.
+ If interval is 5s, link is up for 5s and link is down for 5s
+ """
+
+ while not self.stop_event.is_set():
+ self.stop_event.wait(timeout=self.interval)
+ # simulate link down
+ try:
+ self.packet_drop()
+ log.info('link down')
+ except SSHException as e:
+ log.debug('Failed to run command')
+
+ self.stop_event.wait(timeout=self.interval)
+ # if qdisc exist,delete it.
+ try:
+ delete_dev(self.remote, self.interface)
+ log.info('link up')
+ except SSHException as e:
+ log.debug('Failed to run command')
+
+ def begin(self, gname):
self.thread = gevent.spawn(self.link_toggle)
+ self.ctx.netem.names[gname] = self.thread
- def end(self):
+ def end(self, gname):
+ self.stop_event.set()
+ log.info('gname is {}'.format(self.ctx.netem.names[gname]))
+ self.ctx.netem.names[gname].get()
+
+ def cleanup(self):
+ """
+ Invoked during unwinding if the test fails or exits before executing task 'link_recover'
+ """
+ log.info('Clean up')
self.stop_event.set()
self.thread.get()
- netem:
clients: [rgw.1, mon.0]
iface: eno1
+ gname: t1
dst_client: [c2.rgw.1]
link_toggle_interval: 10 # no unit mentioned. By default takes seconds.
- netem:
clients: [rgw.1, mon.0]
iface: eno1
- link_recover: true
+ link_recover: [t1, t2]
"""
"please list clients to run on"
if not hasattr(ctx, 'netem'):
ctx.netem = argparse.Namespace()
+ ctx.netem.names = {}
if config.get('dst_client') is not None:
dst = config.get('dst_client')
(host,) = ctx.cluster.only(dst).remotes.iterkeys()
for role in config.get('clients', None):
- (remote,) = ctx.cluster.only(role).remotes.iterkeys()
+ (remote,) = ctx.cluster.only(role).remotes.iterkeys()
ctx.netem.remote = remote
if config.get('delay', False):
static_delay(remote, host, config.get('iface'), config.get('delay'))
if config.get('link_toggle_interval', False):
log.info('Toggling link for %s' % config.get('link_toggle_interval'))
global toggle
- toggle = Toggle(remote, host, config.get('iface'), config.get('link_toggle_interval'))
- ctx.netem.toggle = toggle
- toggle.begin()
+ toggle = Toggle(ctx, remote, host, config.get('iface'), config.get('link_toggle_interval'))
+ toggle.begin(config.get('gname'))
if config.get('link_recover', False):
log.info('Recovering link')
- toggle.end()
- log.info('sleeping')
- time.sleep(config.get('link_toggle_interval'))
- delete_dev(ctx.netem.remote, config.get('iface'))
-
+ for gname in config.get('link_recover'):
+ toggle.end(gname)
+ log.info('sleeping')
+ time.sleep(config.get('link_toggle_interval'))
+ delete_dev(ctx.netem.remote, config.get('iface'))
+ del ctx.netem.names[gname]
try:
yield
finally:
- if config.get('link_toggle_interval') and not config.get('link_recover'):
- ctx.netem.toggle.end()
+ if ctx.netem.names:
+ toggle.cleanup()
for role in config.get('clients'):
(remote,) = ctx.cluster.only(role).remotes.iterkeys()
delete_dev(remote, config.get('iface'))