jenkins-bot has submitted this change and it was merged.

Change subject: conftool: add write-locks to syncer and confctl
......................................................................


conftool: add write-locks to syncer and confctl

With this commit, a global distributed lock is acquired when writing to
a specific service or pool, so that concurrent writes are supported.

Also, add a couple more integration tests for confctl.

Bug: T107286
Change-Id: Ib38141256db372cd9becc2e9114ef37bd3087707
---
M conftool/__init__.py
M conftool/cli/syncer.py
M conftool/cli/tool.py
M conftool/drivers/__init__.py
M conftool/drivers/etcd.py
M conftool/tests/integration/test_tool.py
M setup.py
7 files changed, 124 insertions(+), 42 deletions(-)

Approvals:
  Giuseppe Lavagetto: Looks good to me, approved
  jenkins-bot: Verified



diff --git a/conftool/__init__.py b/conftool/__init__.py
index afc68aa..34ec4cd 100644
--- a/conftool/__init__.py
+++ b/conftool/__init__.py
@@ -1,6 +1,7 @@
 import os
 import logging
 import json
+from contextlib import contextmanager
 
 _log = logging.getLogger(__name__)
 from conftool import backend
@@ -66,6 +67,18 @@
             self._set_value(k, self._schema[k], {k: v}, set_defaults=False)
         self.write()
 
+    @classmethod
+    @contextmanager
+    def lock(cls, path):
+        try:
+            l = cls.backend.driver.get_lock(path)
+            yield l
+            cls.backend.driver.release_lock(path)
+        except Exception as e:
+            _log.critical("Problems inside lock for %s: %s", path, e)
+            cls.backend.driver.release_lock(path)
+            raise
+
     def _from_net(self, values):
         """
         Fetch the values from the kvstore into the object
diff --git a/conftool/cli/syncer.py b/conftool/cli/syncer.py
index 20f46de..e268886 100644
--- a/conftool/cli/syncer.py
+++ b/conftool/cli/syncer.py
@@ -59,20 +59,24 @@
 
 def load_services(cluster, servnames, data):
     # TODO: logs, exceptions
-    for servname in servnames:
-        print "Creating service %s/%s" % (cluster, servname)
-        servdata = data[servname]
-        load_service(cluster, servname, servdata)
+    path = service.Service.dir(cluster)
+    with service.Service.lock(path):
+        for servname in servnames:
+            print "Creating service %s/%s" % (cluster, servname)
+            servdata = data[servname]
+            load_service(cluster, servname, servdata)
 
 
 @catch_and_log("error while deleting services")
 def remove_services(cluster, servnames):
-    for servname in servnames:
-        s = service.Service(cluster, servname)
-        if s.exists:
-            print "Removing service %s/%s" % (cluster, servname)
-            _log.info("Removing service %s/%s", cluster, servname)
-            s.delete()
+    path = service.Service.dir(cluster)
+    with service.Service.lock(path):
+        for servname in servnames:
+            s = service.Service(cluster, servname)
+            if s.exists:
+                print "Removing service %s/%s" % (cluster, servname)
+                _log.info("Removing service %s/%s", cluster, servname)
+                s.delete()
 
 
 @catch_and_log("error while calculating changed nodes")
@@ -125,12 +129,14 @@
         for servname, hosts in cl.items():
             new_nodes, del_nodes = get_changed_nodes(dc, cluster,
                                                      servname, hosts)
-            for el in new_nodes:
-                _log.debug("See if %s is present", el)
-                load_node(dc, cluster, servname, el)
-            for el in del_nodes:
-                _log.debug("See if %s should be deleted", el)
-                delete_node(dc, cluster, servname, el)
+            path = node.Node.dir(dc, cluster, servname)
+            with node.Node.lock(path):
+                for el in new_nodes:
+                    _log.debug("See if %s is present", el)
+                    load_node(dc, cluster, servname, el)
+                for el in del_nodes:
+                    _log.debug("See if %s should be deleted", el)
+                    delete_node(dc, cluster, servname, el)
 
 
 def tag_files(directory):
diff --git a/conftool/cli/tool.py b/conftool/cli/tool.py
index 906667d..f20e54b 100644
--- a/conftool/cli/tool.py
+++ b/conftool/cli/tool.py
@@ -107,22 +107,27 @@
     for unit in args.action:
         act, n = unit
         cur_dir = cls.dir(*tags)
-        for name in host_list(n, cur_dir, act):
-            try:
-                # Oh python I <3 you...
-                arguments = list(tags)
-                arguments.append(name)
-                obj = cls(*arguments)
-                a = action.Action(obj, act)
-                msg = a.run()
-            except action.ActionError as e:
-                _log.error("Invalid action, reason: %s", str(e))
-            except BackendError as e:
-                _log.error("Failure writing to the kvstore: %s", str(e))
-            except Exception as e:
-                _log.error("Generic action failure: %s", str(e))
-            else:
-                print(msg)
+        with cls.lock(cur_dir):
+            run_action(cls, n, cur_dir, act, tags)
+
+
+def run_action(cls, n, cur_dir, act, tags):
+    for name in host_list(n, cur_dir, act):
+        try:
+            # Oh python I <3 you...
+            arguments = list(tags)
+            arguments.append(name)
+            obj = cls(*arguments)
+            a = action.Action(obj, act)
+            msg = a.run()
+        except action.ActionError as e:
+            _log.error("Invalid action, reason: %s", str(e))
+        except BackendError as e:
+            _log.error("Failure writing to the kvstore: %s", str(e))
+        except Exception as e:
+            _log.error("Generic action failure: %s", str(e))
+        else:
+            print(msg)
 
 
 if __name__ == '__main__':
diff --git a/conftool/drivers/__init__.py b/conftool/drivers/__init__.py
index 479fc2e..6519897 100644
--- a/conftool/drivers/__init__.py
+++ b/conftool/drivers/__init__.py
@@ -37,6 +37,15 @@
             raise ValueError(
                 "{} is not a directory".format(self.abspath(path)))
 
+    def get_lock(self, path):
+        pass
+
+    def lock_exists(self, path):
+        pass
+
+    def release_lock(self, path):
+        pass
+
 
 def wrap_exception(exc):
     def actual_wrapper(fn):
diff --git a/conftool/drivers/etcd.py b/conftool/drivers/etcd.py
index 9f7e5fb..62e71c5 100644
--- a/conftool/drivers/etcd.py
+++ b/conftool/drivers/etcd.py
@@ -5,10 +5,12 @@
 
 
 class Driver(drivers.BaseDriver):
+    lock_ttl = 60
 
     def __init__(self, config):
         super(Driver, self).__init__(config)
         host_list = []
+        self.locks = {}
         for el in config.hosts:
             h, p = urlparse.urlparse(el).netloc.split(':')
             host_list.append((h, int(p)))
@@ -75,3 +77,31 @@
         if etcdresult is None or etcdresult.dir:
             return None
         return json.loads(etcdresult.value)
+
+    def get_lock(self, path):
+        name = path.replace('/', '-')
+        if name not in self.locks:
+            self.locks[name] = etcd.Lock(self.client, name)
+        self.locks[name].acquire(lock_ttl=self.lock_ttl)
+        if self.locks[name].is_acquired:
+            return self.locks[name]
+        else:
+            return False
+
+    def release_lock(self, path):
+        name = path.replace('/', '-')
+        # we can't remove a lock that was not set by us
+        if name not in self.locks:
+            return False
+        self.locks[name].release()
+        del self.locks[name]
+        return True
+
+    def watch_lock(self, path):
+        name = path.replace('/', '-')
+        l = etcd.Lock(self.client, name)
+        try:
+            r = self.client.read(l.path)
+            return bool(r._children)
+        except etcd.EtcdKeyNotFoundErrror:
+            return False
diff --git a/conftool/tests/integration/test_tool.py 
b/conftool/tests/integration/test_tool.py
index f9c840f..ff69430 100644
--- a/conftool/tests/integration/test_tool.py
+++ b/conftool/tests/integration/test_tool.py
@@ -2,6 +2,7 @@
 import sys
 from conftool.cli import syncer, tool
 from conftool.tests.integration import IntegrationTestBase, test_base
+from conftool import node
 from contextlib import contextmanager
 from StringIO import StringIO
 import json
@@ -9,14 +10,13 @@
 
 @contextmanager
 def captured_output():
-        new_out, new_err = StringIO(), StringIO()
-        old_out, old_err = sys.stdout, sys.stderr
-        try:
-
-            sys.stdout, sys.stderr = new_out, new_err
-            yield sys.stdout, sys.stderr
-        finally:
-            sys.stdout, sys.stderr = old_out, old_err
+    new_out, new_err = StringIO(), StringIO()
+    old_out, old_err = sys.stdout, sys.stderr
+    try:
+        sys.stdout, sys.stderr = new_out, new_err
+        yield sys.stdout, sys.stderr
+    finally:
+        sys.stdout, sys.stderr = old_out, old_err
 
 
 class ToolIntegration(IntegrationTestBase):
@@ -34,10 +34,29 @@
 
     def test_get_node(self):
         args = self.generate_args('get cp1008')
-        print args
         with captured_output() as (out, err):
             tool.main(cmdline=args)
         res = out.getvalue().strip()
         output = json.loads(res)
         self.assertEquals(output.keys(), ['cp1008'])
         self.assertEquals(output['cp1008']['pooled'], 'no')
+
+    def test_change_node_regexp(self):
+        """
+        Test changing values according to a regexp
+        """
+        args = self.generate_args('set/pooled=yes re:cp105.')
+        tool.main(cmdline=args)
+        for hostname in ['cp1052', 'cp1053', 'cp1054', 'cp1055']:
+                n = node.Node('eqiad', 'cache_text', 'https', hostname)
+                self.assertTrue(n.exists)
+                self.assertEquals(n.pooled, "yes")
+
+    def test_create_returns_error(self):
+        """
+        Test creation is not possible from confctl
+        """
+        args = self.generate_args('set/pooled=yes re:cp1059')
+        tool.main(cmdline=args)
+        n = node.Node('eqiad', 'cache_text', 'https', 'cp1059')
+        self.assertFalse(n.exists)
diff --git a/setup.py b/setup.py
index 366a8e7..d114ee2 100755
--- a/setup.py
+++ b/setup.py
@@ -9,7 +9,7 @@
     author='Joe',
     author_email='[email protected]',
     url='https://github.com/wikimedia/operations-software-conftool',
-    install_requires=['python-etcd', 'pyyaml'],
+    install_requires=['python-etcd>=0.4.0', 'pyyaml'],
     test_suite='nose.collector',
     tests_require=['mock', 'nose'],
     zip_safe=True,

-- 
To view, visit https://gerrit.wikimedia.org/r/228015
To unsubscribe, visit https://gerrit.wikimedia.org/r/settings

Gerrit-MessageType: merged
Gerrit-Change-Id: Ib38141256db372cd9becc2e9114ef37bd3087707
Gerrit-PatchSet: 2
Gerrit-Project: operations/software/conftool
Gerrit-Branch: master
Gerrit-Owner: Giuseppe Lavagetto <[email protected]>
Gerrit-Reviewer: Giuseppe Lavagetto <[email protected]>
Gerrit-Reviewer: jenkins-bot <>

_______________________________________________
MediaWiki-commits mailing list
[email protected]
https://lists.wikimedia.org/mailman/listinfo/mediawiki-commits

Reply via email to