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