From 17fac582817e35245bd72395c020620abf8295d2 Mon Sep 17 00:00:00 2001 From: Mike Bonnet Date: Nov 04 2016 18:59:51 +0000 Subject: [PATCH 1/2] protonmsg hub plugin This plugin sends messages to a broker about events in the hub using the proton library. This library supports the AMQP 1.0 protocol and is compatible with a wide variety of message brokers. It queues all messages until the postCommit callback, avoiding race conditions between message reception and database transaction commit. --- diff --git a/plugins/hub/protonmsg.conf b/plugins/hub/protonmsg.conf new file mode 100644 index 0000000..b664196 --- /dev/null +++ b/plugins/hub/protonmsg.conf @@ -0,0 +1,7 @@ +[broker] +urls = amqps://broker1.example.com:5671 amqps://broker2.example.com:5671 +cert = /etc/koji-hub/plugins/client.pem +cacert = /etc/koji-hub/plugins/ca.pem +topic_prefix = koji +connect_timeout = 10 +send_timeout = 60 diff --git a/plugins/hub/protonmsg.py b/plugins/hub/protonmsg.py new file mode 100644 index 0000000..cf3a82c --- /dev/null +++ b/plugins/hub/protonmsg.py @@ -0,0 +1,264 @@ +# Koji callback for sending notifications about events using the +# qpid proton library. +# Copyright (c) 2016 Red Hat, Inc. +# +# Authors: +# Mike Bonnet + +import koji +from koji.plugin import callback, ignore_error +from koji.context import context +import ConfigParser +import logging +import json +import random +from proton import Message, SSLDomain +from proton.reactor import Container +from proton.handlers import MessagingHandler + +CONFIG_FILE = '/etc/koji-hub/plugins/protonmsg.conf' +CONFIG = None + +class TimeoutHandler(MessagingHandler): + def __init__(self, url, msgs, conf, *args, **kws): + super(TimeoutHandler, self).__init__(*args, **kws) + self.url = url + self.msgs = msgs + self.conf = conf + self.pending = {} + self.senders = {} + self.connect_task = None + self.timeout_task = None + self.log = logging.getLogger('koji.plugin.protonmsg.TimeoutHandler') + + def on_start(self, event): + self.log.debug('Container starting') + event.container.connected = False + if self.conf.has_option('broker', 'cert') and self.conf.has_option('broker', 'cacert'): + ssl = SSLDomain(SSLDomain.MODE_CLIENT) + cert = self.conf.get('broker', 'cert') + ssl.set_credentials(cert, cert, None) + ssl.set_trusted_ca_db(self.conf.get('broker', 'cacert')) + ssl.set_peer_authentication(SSLDomain.VERIFY_PEER) + else: + ssl = None + self.log.debug('connecting to %s', self.url) + event.container.connect(url=self.url, reconnect=False, ssl_domain=ssl) + connect_timeout = self.conf.getint('broker', 'connect_timeout') + self.connect_task = event.container.schedule(connect_timeout, self) + send_timeout = self.conf.getint('broker', 'send_timeout') + self.timeout_task = event.container.schedule(send_timeout, self) + + def on_timer_task(self, event): + if not event.container.connected: + self.log.error('not connected, stopping container') + if self.timeout_task: + self.timeout_task.cancel() + self.timeout_task = None + event.container.stop() + else: + # This should only run when called from the timeout task + self.log.error('send timeout expired with %s messages unsent, stopping container', + len(self.msgs)) + event.container.stop() + + def on_connection_opened(self, event): + event.container.connected = True + self.connect_task.cancel() + self.connect_task = None + self.log.debug('connection to %s opened successfully', event.connection.hostname) + self.send_msgs(event) + + def send_msgs(self, event): + prefix = self.conf.get('broker', 'topic_prefix') + for msg in self.msgs: + address = 'topic://' + prefix + '.' + msg[0] + if address in self.senders: + sender = self.senders[address] + self.log.debug('retrieved cached sender for %s', address) + else: + sender = event.container.create_sender(event.connection, target=address) + self.log.debug('created new sender for %s', address) + self.senders[address] = sender + pmsg = Message(properties=msg[1], body=msg[2]) + delivery = sender.send(pmsg) + self.log.debug('sent message: %s', msg[1]) + self.pending[delivery] = msg + + def update_pending(self, event): + msg = self.pending[event.delivery] + del self.pending[event.delivery] + self.log.debug('removed message from self.pending: %s', msg[1]) + if not self.pending: + if self.msgs: + self.log.error('%s messages unsent (rejected or released)', len(self.msgs)) + else: + self.log.debug('all messages sent successfully') + for sender in self.senders.values(): + self.log.debug('closing sender for %s', sender.target.address) + sender.close() + if self.timeout_task: + self.log.debug('canceling timeout task') + self.timeout_task.cancel() + self.timeout_task = None + self.log.debug('closing connection to %s', event.connection.hostname) + event.connection.close() + + def on_settled(self, event): + msg = self.pending[event.delivery] + self.msgs.remove(msg) + self.log.debug('removed message from self.msgs: %s', msg[1]) + self.update_pending(event) + + def on_rejected(self, event): + msg = self.pending[event.delivery] + self.log.error('message was rejected: %s', msg[1]) + self.update_pending(event) + + def on_released(self, event): + msg = self.pending[event.delivery] + self.log.error('message was released: %s', msg[1]) + self.update_pending(event) + + def on_transport_tail_closed(self, event): + if self.connect_task: + self.log.debug('canceling connect timer') + self.connect_task.cancel() + self.connect_task = None + if self.timeout_task: + self.log.debug('canceling send timer') + self.timeout_task.cancel() + self.timeout_task = None + +def queue_msg(address, props, data): + msgs = getattr(context, 'protonmsg_msgs', None) + if msgs is None: + msgs = [] + context.protonmsg_msgs = msgs + body = json.dumps(data) + msgs.append((address, props, body)) + +@callback('postPackageListChange') +def prep_package_list_change(cbtype, *args, **kws): + address = 'package.' + kws['action'] + props = {'type': cbtype[4:], + 'tag': kws['tag']['name'], + 'package': kws['package']['name'], + 'action': kws['action']} + queue_msg(address, props, kws) + +@callback('postTaskStateChange') +def prep_task_state_change(cbtype, *args, **kws): + if kws['attribute'] != 'state': + return + address = 'task.' + kws['new'].lower() + props = {'type': cbtype[4:], + 'id': kws['info']['id'], + 'parent': kws['info']['parent'], + 'method': kws['info']['method'], + 'attribute': kws['attribute'], + 'old': kws['old'], + 'new': kws['new']} + queue_msg(address, props, kws) + +@callback('postBuildStateChange') +def prep_build_state_change(cbtype, *args, **kws): + if kws['attribute'] != 'state': + return + old = kws['old'] + if old is not None: + old = koji.BUILD_STATES[old] + new = koji.BUILD_STATES[kws['new']] + address = 'build.' + new.lower() + props = {'type': cbtype[4:], + 'name': kws['info']['name'], + 'version': kws['info']['version'], + 'release': kws['info']['release'], + 'attribute': kws['attribute'], + 'old': old, + 'new': new} + queue_msg(address, props, kws) + +@callback('postImport') +def prep_import(cbtype, *args, **kws): + address = 'import.' + kws['type'] + props = {'type': cbtype[4:], + 'importType': kws['type'], + 'name': kws['build']['name'], + 'version': kws['build']['version'], + 'release': kws['build']['release']} + queue_msg(address, props, kws) + +@callback('postRPMSign') +def prep_rpm_sign(cbtype, *args, **kws): + address = 'sign.rpm' + props = {'type': cbtype[4:], + 'sigkey': kws['sigkey'], + 'name': kws['build']['name'], + 'version': kws['build']['version'], + 'release': kws['build']['release'], + 'rpm_name': kws['rpm']['name'], + 'rpm_version': kws['rpm']['version'], + 'rpm_release': kws['rpm']['release']} + queue_msg(address, props, kws) + +def _prep_tag_msg(address, cbtype, kws): + build = kws['build'] + props = {'type': cbtype[4:], + 'tag': kws['tag']['name'], + 'name': build['name'], + 'version': build['version'], + 'release': build['release'], + 'user': kws['user']['name']} + queue_msg(address, props, kws) + +@callback('postTag') +def prep_tag(cbtype, *args, **kws): + _prep_tag_msg('build.tag', cbtype, kws) + +@callback('postUntag') +def prep_untag(cbtype, *args, **kws): + _prep_tag_msg('build.untag', cbtype, kws) + +@callback('postRepoInit') +def prep_repo_init(cbtype, *args, **kws): + address = 'repo.init' + props = {'type': cbtype[4:], + 'tag': kws['tag']['name'], + 'repo_id': kws['repo_id']} + queue_msg(address, props, kws) + +@callback('postRepoDone') +def prep_repo_done(cbtype, *args, **kws): + address = 'repo.done' + props = {'type': cbtype[4:], + 'tag': kws['repo']['tag_name'], + 'repo_id': kws['repo']['id'], + 'expire': kws['expire']} + queue_msg(address, props, kws) + +@ignore_error +@callback('postCommit') +def send_queued_msgs(cbtype, *args, **kws): + msgs = getattr(context, 'protonmsg_msgs', None) + if not msgs: + return + log = logging.getLogger('koji.plugin.protonmsg') + global CONFIG + if not CONFIG: + conf = ConfigParser.SafeConfigParser() + with open(CONFIG_FILE) as conffile: + conf.readfp(conffile) + CONFIG = conf + urls = CONFIG.get('broker', 'urls').split() + for url in sorted(urls, key=lambda k: random.random()): + container = Container(TimeoutHandler(url, msgs, CONFIG)) + container.run() + if msgs: + log.debug('could not send to %s, %s messages remaining', + url, len(msgs)) + else: + log.debug('all messages sent to %s successfully', url) + break + else: + log.error('could not send messages to any destinations') diff --git a/tests/test_plugins/test_protonmsg.py b/tests/test_plugins/test_protonmsg.py new file mode 100644 index 0000000..02583b7 --- /dev/null +++ b/tests/test_plugins/test_protonmsg.py @@ -0,0 +1,365 @@ +import unittest +from mock import patch, MagicMock +import protonmsg +from koji.context import context +import tempfile +from StringIO import StringIO +from ConfigParser import SafeConfigParser + +class TestProtonMsg(unittest.TestCase): + def tearDown(self): + if hasattr(context, 'protonmsg_msgs'): + del context.protonmsg_msgs + + def assertMsg(self, topic, body=None, **kws): + self.assertTrue(hasattr(context, 'protonmsg_msgs')) + self.assertEqual(len(context.protonmsg_msgs), 1) + msg = context.protonmsg_msgs[0] + self.assertEqual(msg[0], topic) + for kw in kws: + self.assertTrue(kw in msg[1]) + self.assertEqual(msg[1][kw], kws[kw]) + self.assertEqual(len(msg[1]), len(kws)) + if body: + self.assertEqual(msg[2], body) + + def test_queue_msg(self): + protonmsg.queue_msg('test.msg', {'testheader': 1}, 'test body') + self.assertMsg('test.msg', body='"test body"', testheader=1) + + def test_prep_package_list_change_add(self): + protonmsg.prep_package_list_change('postPackageListChange', + action='add', tag={'name': 'test-tag'}, + package={'name': 'test-pkg'}, + owner=1, + block=False, extra_arches='i386 x86_64', + force=False, update=False) + self.assertMsg('package.add', type='PackageListChange', tag='test-tag', + package='test-pkg', action='add') + + def test_prep_package_list_change_update(self): + protonmsg.prep_package_list_change('postPackageListChange', + action='update', tag={'name': 'test-tag'}, + package={'name': 'test-pkg'}, + owner=1, + block=False, extra_arches='i386 x86_64', + force=False, update=False) + self.assertMsg('package.update', type='PackageListChange', tag='test-tag', + package='test-pkg', action='update') + + def test_prep_package_list_change_block(self): + protonmsg.prep_package_list_change('postPackageListChange', + action='block', tag={'name': 'test-tag'}, + package={'name': 'test-pkg'}, + owner=1, + block=False, extra_arches='i386 x86_64', + force=False, update=False) + self.assertMsg('package.block', type='PackageListChange', tag='test-tag', + package='test-pkg', action='block') + + def test_prep_package_list_change_unblock(self): + protonmsg.prep_package_list_change('postPackageListChange', + action='unblock', tag={'name': 'test-tag'}, + package={'name': 'test-pkg'}) + self.assertMsg('package.unblock', type='PackageListChange', tag='test-tag', + package='test-pkg', action='unblock') + + def test_prep_package_list_change_remove(self): + protonmsg.prep_package_list_change('postPackageListChange', + action='remove', tag={'name': 'test-tag'}, + package={'name': 'test-pkg'}) + self.assertMsg('package.remove', type='PackageListChange', tag='test-tag', + package='test-pkg', action='remove') + + def test_prep_task_state_change(self): + info = {'id': 5678, + 'parent': 1234, + 'method': 'build'} + protonmsg.prep_task_state_change('postTaskStateChange', + info=info, attribute='weight', + old=2.0, new=3.5) + # no messages should be created for callbacks where attribute != state + self.assertFalse(hasattr(context, 'protonmsg_msgs')) + protonmsg.prep_task_state_change('postTaskStateChange', + info=info, attribute='state', + old='FREE', new='OPEN') + self.assertMsg('task.open', type='TaskStateChange', + attribute='state', old='FREE', new='OPEN', + **info) + + def test_prep_build_state_change(self): + info = {'name': 'test-pkg', + 'version': '1.0', + 'release': '1'} + protonmsg.prep_build_state_change('postBuildStateChange', + info=info, attribute='volume_id', + old=0, new=1) + # no messages should be created for callbacks where attribute != state + self.assertFalse(hasattr(context, 'protonmsg_msgs')) + protonmsg.prep_build_state_change('postBuildStateChange', + info=info, attribute='state', + old=0, new=1) + self.assertMsg('build.complete', type='BuildStateChange', + attribute='state', old='BUILDING', new='COMPLETE', + **info) + + def test_prep_import(self): + build = {'name': 'test-pkg', 'version': '1.0', 'release': '1'} + protonmsg.prep_import('postImport', type='build', build=build) + self.assertMsg('import.build', type='Import', importType='build', + **build) + + def test_prep_rpm_sign(self): + build = {'name': 'test-pkg', + 'version': '1.0', + 'release': '1'} + rpm = {'name': 'test-pkg-subpkg', + 'version': '2.0', + 'release': '2'} + sigkey = 'a1b2c3d4' + protonmsg.prep_rpm_sign('postRPMSign', sigkey=sigkey, sighash='fedcba9876543210', + build=build, rpm=rpm) + self.assertMsg('sign.rpm', type='RPMSign', sigkey=sigkey, rpm_name=rpm['name'], + rpm_version=rpm['version'], rpm_release=rpm['release'], + **build) + + def test_prep_tag(self): + build = {'name': 'test-pkg', 'version': '1.0', 'release': '1'} + protonmsg.prep_tag('postTag', tag={'name': 'test-tag'}, + build=build, user={'name': 'test-user'}) + self.assertMsg('build.tag', type='Tag', tag='test-tag', + user='test-user', **build) + + def test_prep_untag(self): + build = {'name': 'test-pkg', 'version': '1.0', 'release': '1'} + protonmsg.prep_untag('postUntag', tag={'name': 'test-tag'}, + build=build, user={'name': 'test-user'}) + self.assertMsg('build.untag', type='Untag', tag='test-tag', + user='test-user', **build) + + def test_prep_repo_init(self): + protonmsg.prep_repo_init('postRepoInit', tag={'name': 'test-tag'}, repo_id=1234) + self.assertMsg('repo.init', type='RepoInit', tag='test-tag', repo_id=1234) + + def test_prep_repo_done(self): + protonmsg.prep_repo_done('postRepoDone', repo={'tag_name': 'test-tag', 'id': 1234}, + expire=False) + self.assertMsg('repo.done', type='RepoDone', tag='test-tag', repo_id=1234, expire=False) + + @patch('protonmsg.Container') + def test_send_queued_msgs_none(self, Container): + self.assertFalse(hasattr(context, 'protonmsg_msgs')) + protonmsg.send_queued_msgs('postCommit') + self.assertEqual(Container.call_count, 0) + context.protonmsg_msgs = [] + protonmsg.send_queued_msgs('postCommit') + self.assertEqual(Container.call_count, 0) + + @patch('protonmsg.Container') + @patch('logging.getLogger') + def test_send_queued_msgs_fail(self, getLogger, Container): + context.protonmsg_msgs = [('test.topic', {'testheader': 1}, 'test body')] + conf = tempfile.NamedTemporaryFile() + conf.write("""[broker] +urls = amqps://broker1.example.com:5671 amqps://broker2.example.com:5671 +cert = /etc/koji-hub/plugins/client.pem +cacert = /etc/koji-hub/plugins/ca.pem +topic_prefix = koji +connect_timeout = 10 +send_timeout = 60 +""") + conf.flush() + protonmsg.CONFIG_FILE = conf.name + protonmsg.send_queued_msgs('postCommit') + log = getLogger.return_value + self.assertEqual(log.debug.call_count, 2) + for args in log.debug.call_args_list: + self.assertTrue(args[0][0].startswith('could not send')) + self.assertEqual(log.error.call_count, 1) + self.assertTrue(log.error.call_args[0][0].startswith('could not send')) + + @patch('protonmsg.Container') + @patch('logging.getLogger') + def test_send_queued_msgs_success(self, getLogger, Container): + context.protonmsg_msgs = [('test.topic', {'testheader': 1}, 'test body')] + conf = tempfile.NamedTemporaryFile() + conf.write("""[broker] +urls = amqps://broker1.example.com:5671 amqps://broker2.example.com:5671 +cert = /etc/koji-hub/plugins/client.pem +cacert = /etc/koji-hub/plugins/ca.pem +topic_prefix = koji +connect_timeout = 10 +send_timeout = 60 +""") + conf.flush() + protonmsg.CONFIG_FILE = conf.name + def clear_msgs(): + del context.protonmsg_msgs[:] + Container.return_value.run.side_effect = clear_msgs + protonmsg.send_queued_msgs('postCommit') + log = getLogger.return_value + self.assertEqual(log.debug.call_count, 1) + self.assertTrue(log.debug.args[0][0].startswith('all msgs sent')) + self.assertEqual(log.error.call_count, 0) + +class TestTimeoutHandler(unittest.TestCase): + def setUp(self): + confdata = StringIO("""[broker] +urls = amqps://broker1.example.com:5671 amqps://broker2.example.com:5671 +cert = /etc/koji-hub/plugins/client.pem +cacert = /etc/koji-hub/plugins/ca.pem +topic_prefix = koji +connect_timeout = 10 +send_timeout = 60 +""") + conf = SafeConfigParser() + conf.readfp(confdata) + self.handler = protonmsg.TimeoutHandler('amqps://broker1.example.com:5671', [], conf) + + @patch('protonmsg.SSLDomain') + def test_on_start(self, SSLDomain): + event = MagicMock() + self.handler.on_start(event) + event.container.connect.assert_called_once_with(url='amqps://broker1.example.com:5671', + reconnect=False, + ssl_domain=SSLDomain.return_value) + self.assertEqual(event.container.schedule.call_count, 2) + + @patch('protonmsg.SSLDomain') + def test_on_start_no_ssl(self, SSLDomain): + confdata = StringIO("""[broker] +urls = amqp://broker1.example.com:5672 amqp://broker2.example.com:5672 +topic_prefix = koji +connect_timeout = 10 +send_timeout = 60 +""") + conf = SafeConfigParser() + conf.readfp(confdata) + handler = protonmsg.TimeoutHandler('amqp://broker1.example.com:5672', [], conf) + event = MagicMock() + handler.on_start(event) + event.container.connect.assert_called_once_with(url='amqp://broker1.example.com:5672', + reconnect=False, + ssl_domain=None) + self.assertEqual(SSLDomain.call_count, 0) + + @patch('protonmsg.SSLDomain') + def test_on_timer_task(self, SSLDomain): + event = MagicMock() + self.handler.on_start(event) + self.assertTrue(self.handler.timeout_task is not None) + self.handler.on_timer_task(event) + event.container.schedule.return_value.cancel.assert_called_once_with() + self.assertTrue(self.handler.timeout_task is None) + event.container.stop.assert_called_once_with() + event.container.stop.reset_mock() + self.handler.log = MagicMock() + event.container.connected = True + self.handler.on_timer_task(event) + event.container.stop.assert_called_once_with() + self.assertTrue(self.handler.log.error.call_args[0][0].startswith('send timeout expired')) + + @patch('protonmsg.SSLDomain') + def test_on_connection_opened(self, SSLDomain): + event = MagicMock() + self.handler.on_start(event) + self.assertTrue(self.handler.connect_task is not None) + self.handler.on_connection_opened(event) + self.assertTrue(event.container.connected) + event.container.schedule.return_value.cancel.assert_called_once_with() + self.assertTrue(self.handler.connect_task is None) + + @patch('protonmsg.Message') + @patch('protonmsg.SSLDomain') + def test_send_msgs(self, SSLDomain, Message): + event = MagicMock() + self.handler.on_start(event) + self.handler.msgs = [('testtopic', {'testheader': 1}, '"test body"')] + self.handler.on_connection_opened(event) + event.container.create_sender.assert_called_once_with(event.connection, + target='topic://koji.testtopic') + Message.assert_called_once_with(properties={'testheader': 1}, body='"test body"') + sender = event.container.create_sender.return_value + sender.send.assert_called_once_with(Message.return_value) + + @patch('protonmsg.Message') + @patch('protonmsg.SSLDomain') + def test_update_pending(self, SSLDomain, Message): + event = MagicMock() + self.handler.on_start(event) + self.handler.msgs = [('testtopic', {'testheader': 1}, '"test body"'), + ('testtopic', {'testheader': 2}, '"test body"')] + delivery0 = MagicMock() + delivery1 = MagicMock() + sender = event.container.create_sender.return_value + sender.send.side_effect = [delivery0, delivery1] + log = MagicMock() + self.handler.log = log + self.handler.on_connection_opened(event) + self.assertEqual(len(self.handler.pending), 2) + event.delivery = delivery0 + self.handler.update_pending(event) + self.assertEqual(len(self.handler.pending), 1) + self.assertTrue(delivery0 not in self.handler.pending) + log.debug.call_args[0][0].startswith('removed msg') + event.delivery = delivery1 + self.handler.update_pending(event) + self.assertEqual(len(self.handler.pending), 0) + self.assertTrue(delivery0 not in self.handler.pending) + log.error.call_args[0][0].startswith('2 messages unsent') + sender.close.assert_called_once_with() + self.assertEqual(event.container.schedule.return_value.cancel.call_count, 2) + event.connection.close.assert_called_once_with() + + @patch('protonmsg.Message') + @patch('protonmsg.SSLDomain') + def test_on_settled(self, SSLDomain, Message): + event = MagicMock() + self.handler.on_start(event) + self.handler.msgs = [('testtopic', {'testheader': 1}, '"test body"')] + self.handler.on_connection_opened(event) + delivery = event.container.create_sender.return_value.send.return_value + self.assertTrue(delivery in self.handler.pending) + event.delivery = delivery + self.handler.on_settled(event) + self.assertEqual(len(self.handler.msgs), 0) + self.assertEqual(len(self.handler.pending), 0) + + @patch('protonmsg.Message') + @patch('protonmsg.SSLDomain') + def test_on_rejected(self, SSLDomain, Message): + event = MagicMock() + self.handler.on_start(event) + self.handler.msgs = [('testtopic', {'testheader': 1}, '"test body"')] + self.handler.on_connection_opened(event) + delivery = event.container.create_sender.return_value.send.return_value + self.assertTrue(delivery in self.handler.pending) + event.delivery = delivery + self.handler.on_rejected(event) + self.assertEqual(len(self.handler.msgs), 1) + self.assertEqual(len(self.handler.pending), 0) + + @patch('protonmsg.Message') + @patch('protonmsg.SSLDomain') + def test_on_released(self, SSLDomain, Message): + event = MagicMock() + self.handler.on_start(event) + self.handler.msgs = [('testtopic', {'testheader': 1}, '"test body"')] + self.handler.on_connection_opened(event) + delivery = event.container.create_sender.return_value.send.return_value + self.assertTrue(delivery in self.handler.pending) + event.delivery = delivery + self.handler.on_released(event) + self.assertEqual(len(self.handler.msgs), 1) + self.assertEqual(len(self.handler.pending), 0) + + @patch('protonmsg.SSLDomain') + def test_on_transport_tail_closed(self, SSLDomain): + event = MagicMock() + self.handler.on_start(event) + self.assertTrue(self.handler.connect_task is not None) + self.assertTrue(self.handler.timeout_task is not None) + self.handler.on_transport_tail_closed(event) + self.assertEqual(event.container.schedule.return_value.cancel.call_count, 2) + self.assertTrue(self.handler.connect_task is None) + self.assertTrue(self.handler.timeout_task is None) From f09d3115b0c9cced0dd77985ce0c0f2f2a4d0076 Mon Sep 17 00:00:00 2001 From: Mike Bonnet Date: Nov 04 2016 18:59:51 +0000 Subject: [PATCH 2/2] add dependency on python-qpid-proton to koji-hub-plugins python-qpid-proton is required by the protonmsg hub plugin --- diff --git a/koji.spec b/koji.spec index ea758a1..5d12d27 100644 --- a/koji.spec +++ b/koji.spec @@ -66,6 +66,7 @@ License: LGPLv2 Requires: %{name} = %{version}-%{release} Requires: %{name}-hub = %{version}-%{release} Requires: python-qpid >= 0.7 +Requires: python-qpid-proton %if 0%{?rhel} == 5 Requires: python-ssl %endif