events: Move subscriber classes from svn_support to rhodecode.subscribers.

These classes are generic and can be used by other components as well.
Thereofore they should not live inside of the svn_support package.
This commit is contained in:
Martin Bornhold 2016-10-17 10:44:56 +02:00
parent 472f756d3d
commit ec8d1af12f
3 changed files with 75 additions and 70 deletions

View file

@ -19,14 +19,21 @@
# and proprietary license terms, please see https://rhodecode.com/licenses/
import logging
import pylons
import Queue
import subprocess32
from pyramid.i18n import get_localizer
from pyramid.threadlocal import get_current_request
from threading import Thread
from rhodecode.translation import _ as tsf
log = logging.getLogger(__name__)
def add_renderer_globals(event):
# Put pylons stuff into the context. This will be removed as soon as
# migration to pyramid is finished.
@ -68,3 +75,69 @@ def scan_repositories_if_enabled(event):
if vcs_server_enabled and import_on_startup:
repositories = ScmModel().repo_scan(get_rhodecode_base_path())
repo2db_mapper(repositories, remove_obsolete=False)
class Subscriber(object):
def __call__(self, event):
self.run(event)
def run(self, event):
raise NotImplementedError('Subclass has to implement this.')
class AsyncSubscriber(Subscriber):
def __init__(self):
self._stop = False
self._eventq = Queue.Queue()
self._worker = self.create_worker()
self._worker.start()
def __call__(self, event):
self._eventq.put(event)
def create_worker(self):
worker = Thread(target=self.do_work)
worker.daemon = True
return worker
def stop_worker(self):
self._stop = False
self._eventq.put(None)
self._worker.join()
def do_work(self):
while not self._stop:
event = self._eventq.get()
if event is not None:
self.run(event)
class AsyncSubprocessSubscriber(AsyncSubscriber):
def __init__(self, cmd, timeout=None):
super(AsyncSubprocessSubscriber, self).__init__()
self._cmd = cmd
self._timeout = timeout
def run(self, event):
cmd = self._cmd
timeout = self._timeout
log.debug('Executing command %s.', cmd)
try:
output = subprocess32.check_output(
cmd, timeout=timeout, stderr=subprocess32.STDOUT)
log.debug('Command finished %s', cmd)
if output:
log.debug('Command output: %s', output)
except subprocess32.TimeoutExpired as e:
log.exception('Timeout while executing command.')
if e.output:
log.error('Command output: %s', e.output)
except subprocess32.CalledProcessError as e:
log.exception('Error while executing command.')
if e.output:
log.error('Command output: %s', e.output)
except:
log.exception(
'Exception while executing command %s.', cmd)

View file

@ -25,11 +25,12 @@ import shlex
# Do not use `from rhodecode import events` here, it will be overridden by the
# events module in this package due to pythons import mechanism.
from rhodecode.events import RepoGroupEvent
from rhodecode.subscribers import AsyncSubprocessSubscriber
from rhodecode.config.middleware import (
_bool_setting, _string_setting, _int_setting)
from .events import ModDavSvnConfigChange
from .subscribers import generate_config_subscriber, AsyncSubprocessSubscriber
from .subscribers import generate_config_subscriber
from . import config_keys

View file

@ -19,9 +19,6 @@
# and proprietary license terms, please see https://rhodecode.com/licenses/
import logging
import Queue
import subprocess32
from threading import Thread
from .utils import generate_mod_dav_svn_config
@ -41,69 +38,3 @@ def generate_config_subscriber(event):
except Exception:
log.exception(
'Exception while generating subversion mod_dav_svn configuration.')
class Subscriber(object):
def __call__(self, event):
self.run(event)
def run(self, event):
raise NotImplementedError('Subclass has to implement this.')
class AsyncSubscriber(Subscriber):
def __init__(self):
self._stop = False
self._eventq = Queue.Queue()
self._worker = self.create_worker()
self._worker.start()
def __call__(self, event):
self._eventq.put(event)
def create_worker(self):
worker = Thread(target=self.do_work)
worker.daemon = True
return worker
def stop_worker(self):
self._stop = False
self._eventq.put(None)
self._worker.join()
def do_work(self):
while not self._stop:
event = self._eventq.get()
if event is not None:
self.run(event)
class AsyncSubprocessSubscriber(AsyncSubscriber):
def __init__(self, cmd, timeout=None):
super(AsyncSubprocessSubscriber, self).__init__()
self._cmd = cmd
self._timeout = timeout
def run(self, event):
cmd = self._cmd
timeout = self._timeout
log.debug('Executing command %s.', cmd)
try:
output = subprocess32.check_output(
cmd, timeout=timeout, stderr=subprocess32.STDOUT)
log.debug('Command finished %s', cmd)
if output:
log.debug('Command output: %s', output)
except subprocess32.TimeoutExpired as e:
log.exception('Timeout while executing command.')
if e.output:
log.error('Command output: %s', e.output)
except subprocess32.CalledProcessError as e:
log.exception('Error while executing command.')
if e.output:
log.error('Command output: %s', e.output)
except:
log.exception(
'Exception while executing command %s.', cmd)