#
# Author: Alina Quereilhac <alina.quereilhac@inria.fr>
-from nepi.execution.attribute import Attribute, Flags
-from nepi.execution.resource import ResourceManager, clsinit, ResourceState, \
- reschedule_delay
+from nepi.execution.attribute import Attribute, Flags, Types
+from nepi.execution.resource import ResourceManager, clsinit_copy, \
+ ResourceState, reschedule_delay, failtrap
from nepi.resources.linux import rpmfuncs, debfuncs
from nepi.util import sshfuncs, execfuncs
from nepi.util.sshfuncs import ProcStatus
UBUNTU = "ubuntu"
DEBIAN = "debian"
-@clsinit
+@clsinit_copy
class LinuxNode(ResourceManager):
"""
.. class:: Class Args :
"""
_rtype = "LinuxNode"
+ _help = "Controls Linux host machines ( either localhost or a host " \
+ "that can be accessed using a SSH key)"
+ _backend_type = "linux"
@classmethod
def _register_attributes(cls):
clean_home = Attribute("cleanHome", "Remove all nepi files and directories "
" from node home folder before starting experiment",
+ type = Types.Bool,
+ default = False,
flags = Flags.ExecReadOnly)
clean_experiment = Attribute("cleanExperiment", "Remove all files and directories "
" from a previous same experiment, before the new experiment starts",
+ type = Types.Bool,
+ default = False,
flags = Flags.ExecReadOnly)
clean_processes = Attribute("cleanProcesses",
"Kill all running processes before starting experiment",
+ type = Types.Bool,
+ default = False,
flags = Flags.ExecReadOnly)
tear_down = Attribute("tearDown", "Bash script to be executed before " + \
# home directory at Linux host
self._home_dir = ""
- # lock to avoid concurrency issues on methods used by applications
- self._lock = threading.Lock()
+ # lock to prevent concurrent applications on the same node,
+ # to execute commands at the same time. There are potential
+ # concurrency issues when using SSH to a same host from
+ # multiple threads. There are also possible operational
+ # issues, e.g. an application querying the existence
+ # of a file or folder prior to its creation, and another
+ # application creating the same file or folder in between.
+ self._node_lock = threading.Lock()
def log_message(self, msg):
return " guid %d - host %s - %s " % (self.guid,
self.error(msg)
raise RuntimeError, msg
- (out, err), proc = self.execute("cat /etc/issue", with_lock = True)
-
- if err and proc.poll():
- msg = "Error detecting OS "
- self.error(msg, out, err)
- raise RuntimeError, "%s - %s - %s" %( msg, out, err )
+ out = self.get_os()
if out.find("Fedora release 8") == 0:
self._os = OSType.FEDORA_8
return self._os
+ def get_os(self):
+ # The underlying SSH layer will sometimes return an empty
+ # output (even if the command was executed without errors).
+ # To work arround this, repeat the operation N times or
+ # until the result is not empty string
+ out = ""
+ retrydelay = 1.0
+ for i in xrange(10):
+ try:
+ (out, err), proc = self.execute("cat /etc/issue",
+ retry = 5,
+ with_lock = True,
+ blocking = True)
+
+ if out.strip() != "":
+ return out
+ except:
+ trace = traceback.format_exc()
+ msg = "Error detecting OS: %s " % trace
+ self.error(msg, out, err)
+ return False
+
+ time.sleep(min(30.0, retrydelay))
+ retrydelay *= 1.5
+
@property
def use_deb(self):
return self.os in [OSType.DEBIAN, OSType.UBUNTU]
def localhost(self):
return self.get("hostname") in ['localhost', '127.0.0.7', '::1']
+ @failtrap
def provision(self):
# check if host is alive
if not self.is_alive():
- self.fail()
-
msg = "Deploy failed. Unresponsive node %s" % self.get("hostname")
self.error(msg)
raise RuntimeError, msg
super(LinuxNode, self).provision()
+ @failtrap
def deploy(self):
if self.state == ResourceState.NEW:
- try:
- self.discover()
- self.provision()
- except:
- self._state = ResourceState.FAILED
- raise
+ self.info("Deploying node")
+ self.discover()
+ self.provision()
# Node needs to wait until all associated interfaces are
# ready before it can finalize deployment
from nepi.resources.linux.interface import LinuxInterface
- ifaces = self.get_connected(LinuxInterface)
+ ifaces = self.get_connected(LinuxInterface.rtype())
for iface in ifaces:
if iface.state < ResourceState.READY:
self.ec.schedule(reschedule_delay, self.deploy)
super(LinuxNode, self).deploy()
def release(self):
- tear_down = self.get("tearDown")
- if tear_down:
- self.execute(tear_down)
+ try:
+ rms = self.get_connected()
+ for rm in rms:
+ # Node needs to wait until all associated RMs are released
+ # before it can be released
+ if rm.state < ResourceState.STOPPED:
+ self.ec.schedule(reschedule_delay, self.release)
+ return
+
+ tear_down = self.get("tearDown")
+ if tear_down:
+ self.execute(tear_down)
- self.clean_processes()
+ self.clean_processes()
+ except:
+ import traceback
+ err = traceback.format_exc()
+ self.error(err)
super(LinuxNode, self).release()
env = env)
else:
if with_lock:
- with self._lock:
+ with self._node_lock:
(out, err), proc = sshfuncs.rexec(
command,
host = self.get("hostname"),
sudo = sudo,
user = user)
else:
- with self._lock:
+ with self._node_lock:
(out, err), proc = sshfuncs.rspawn(
command,
pidfile = pidfile,
if self.localhost:
pidtuple = execfuncs.lgetpid(os.path.join(home, pidfile))
else:
- with self._lock:
+ with self._node_lock:
pidtuple = sshfuncs.rgetpid(
os.path.join(home, pidfile),
host = self.get("hostname"),
if self.localhost:
status = execfuncs.lstatus(pid, ppid)
else:
- with self._lock:
+ with self._node_lock:
status = sshfuncs.rstatus(
pid, ppid,
host = self.get("hostname"),
if self.localhost:
(out, err), proc = execfuncs.lkill(pid, ppid, sudo)
else:
- with self._lock:
+ with self._node_lock:
(out, err), proc = sshfuncs.rkill(
pid, ppid,
host = self.get("hostname"),
recursive = True,
strict_host_checking = False)
else:
- with self._lock:
+ with self._node_lock:
(out, err), proc = sshfuncs.rcopy(
src, dst,
port = self.get("port"),
return (out, err), proc
-
def upload(self, src, dst, text = False, overwrite = True):
""" Copy content to destination
src = "%s@%s:%s" % (self.get("username"), self.get("hostname"), src)
return self.copy(src, dst)
- def install_packages(self, packages, home, run_home = None):
- """ Install packages in the Linux host.
-
- 'home' is the directory to upload the package installation script.
- 'run_home' is the directory from where to execute the script.
- """
+ def install_packages_command(self, packages):
command = ""
if self.use_rpm:
command = rpmfuncs.install_packages_command(self.os, packages)
self.error(msg, self.os)
raise RuntimeError, msg
+ return command
+
+ def install_packages(self, packages, home, run_home = None):
+ """ Install packages in the Linux host.
+
+ 'home' is the directory to upload the package installation script.
+ 'run_home' is the directory from where to execute the script.
+ """
+ command = self.install_packages_command(packages)
+
run_home = run_home or home
(out, err), proc = self.run_and_wait(command, run_home,
def wait_run(self, pid, ppid, trial = 0):
""" wait for a remote process to finish execution """
- start_delay = 1.0
+ delay = 1.0
while True:
status = self.status(pid, ppid)
return True
out = err = ""
- try:
- (out, err), proc = self.execute("echo 'ALIVE'",
- retry = 5,
- with_lock = True)
- except:
- trace = traceback.format_exc()
- msg = "Unresponsive host %s " % err
- self.error(msg, out, trace)
- return False
+ # The underlying SSH layer will sometimes return an empty
+ # output (even if the command was executed without errors).
+ # To work arround this, repeat the operation N times or
+ # until the result is not empty string
+ retrydelay = 1.0
+ for i in xrange(10):
+ try:
+ (out, err), proc = self.execute("echo 'ALIVE'",
+ retry = 5,
+ blocking = True,
+ with_lock = True)
+
+ if out.find("ALIVE") > -1:
+ return True
+ except:
+ trace = traceback.format_exc()
+ msg = "Unresponsive host. Error reaching host: %s " % trace
+ self.error(msg, out, err)
+ return False
+
+ time.sleep(min(30.0, retrydelay))
+ retrydelay *= 1.5
- if out.strip() == "ALIVE":
+ if out.find("ALIVE") > -1:
return True
else:
- msg = "Unresponsive host "
+ msg = "Unresponsive host. Wrong answer. "
self.error(msg, out, err)
return False
def find_home(self):
""" Retrieves host home directory
"""
- (out, err), proc = self.execute("echo ${HOME}", retry = 5,
- with_lock = True)
+ # The underlying SSH layer will sometimes return an empty
+ # output (even if the command was executed without errors).
+ # To work arround this, repeat the operation N times or
+ # until the result is not empty string
+ retrydelay = 1.0
+ for i in xrange(10):
+ try:
+ (out, err), proc = self.execute("echo ${HOME}",
+ retry = 5,
+ blocking = True,
+ with_lock = True)
+
+ if out.strip() != "":
+ self._home_dir = out.strip()
+ break
+ except:
+ trace = traceback.format_exc()
+ msg = "Impossible to retrieve HOME directory" % trace
+ self.error(msg, out, err)
+ return False
- if proc.poll():
- msg = "Imposible to retrieve HOME directory"
+ time.sleep(min(30.0, retrydelay))
+ retrydelay *= 1.5
+
+ if not self._home_dir:
+ msg = "Impossible to retrieve HOME directory"
self.error(msg, out, err)
raise RuntimeError, msg
- self._home_dir = out.strip()
-
def filter_existing_files(self, src, dst):
""" Removes files that already exist in the Linux host from src list
"""