Skip to content

Commit 9027154

Browse files
committed
Removing dependency on eventlet and oslo.service
Change-Id: I453e9b86d4edfedd63cc59e47bf745e166ff836f
1 parent 84b5078 commit 9027154

7 files changed

Lines changed: 58 additions & 55 deletions

File tree

devstack/plugin.sh

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -134,6 +134,7 @@ function octavia_configure {
134134
iniuncomment $OCTAVIA_CONF haproxy_amphora rest_request_read_timeout
135135
iniuncomment $OCTAVIA_CONF controller_worker amp_active_retries
136136
iniuncomment $OCTAVIA_CONF controller_worker amp_active_wait_sec
137+
iniuncomment $OCTAVIA_CONF controller_worker workers
137138

138139
# devstack optimizations for tempest runs
139140
iniset $OCTAVIA_CONF haproxy_amphora connection_max_retries 1500
@@ -142,6 +143,7 @@ function octavia_configure {
142143
iniset $OCTAVIA_CONF haproxy_amphora rest_request_read_timeout ${OCTAVIA_AMP_READ_TIMEOUT}
143144
iniset $OCTAVIA_CONF controller_worker amp_active_retries 100
144145
iniset $OCTAVIA_CONF controller_worker amp_active_wait_sec 2
146+
iniset $OCTAVIA_CONF controller_worker workers 2
145147

146148
if [[ -a $OCTAVIA_SSH_DIR ]] ; then
147149
rm -rf $OCTAVIA_SSH_DIR

etc/octavia.conf

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -147,6 +147,7 @@
147147
# rest_request_read_timeout = 60
148148

149149
[controller_worker]
150+
# workers = 1
150151
# amp_active_retries = 10
151152
# amp_active_wait_sec = 10
152153
# Glance parameters to extract image ID to use for amphora. Only one of

octavia/cmd/octavia_worker.py

Lines changed: 12 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -14,15 +14,14 @@
1414

1515
import sys
1616

17-
import eventlet
18-
eventlet.monkey_patch()
19-
from oslo_config import cfg # noqa: E402
20-
from oslo_reports import guru_meditation_report as gmr # noqa: E402
21-
from oslo_service import service # noqa: E402
17+
import cotyledon
18+
from cotyledon import oslo_config_glue
19+
from oslo_config import cfg
20+
from oslo_reports import guru_meditation_report as gmr
2221

23-
from octavia.common import service as octavia_service # noqa: E402
24-
from octavia.controller.queue import consumer # noqa: E402
25-
from octavia import version # noqa: E402
22+
from octavia.common import service as octavia_service
23+
from octavia.controller.queue import consumer
24+
from octavia import version
2625

2726
CONF = cfg.CONF
2827

@@ -32,5 +31,8 @@ def main():
3231

3332
gmr.TextGuruMeditation.setup_autorun(version)
3433

35-
launcher = service.launch(CONF, consumer.Consumer())
36-
launcher.wait()
34+
sm = cotyledon.ServiceManager()
35+
sm.add(consumer.ConsumerService, workers=CONF.controller_worker.workers,
36+
args=(CONF,))
37+
oslo_config_glue.setup(sm, CONF)
38+
sm.run()

octavia/common/config.py

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -216,6 +216,9 @@
216216
]
217217

218218
controller_worker_opts = [
219+
cfg.IntOpt('workers',
220+
default=1, min=1,
221+
help='Number of workers for the controller-worker service.'),
219222
cfg.IntOpt('amp_active_retries',
220223
default=10,
221224
help=_('Retry attempts to wait for Amphora to become active')),

octavia/controller/queue/consumer.py

Lines changed: 24 additions & 29 deletions
Original file line numberDiff line numberDiff line change
@@ -12,50 +12,45 @@
1212
# License for the specific language governing permissions and limitations
1313
# under the License.
1414

15-
from oslo_config import cfg
15+
import cotyledon
1616
from oslo_log import log as logging
1717
import oslo_messaging as messaging
1818
from oslo_messaging.rpc import dispatcher
19-
from oslo_service import service
2019

2120
from octavia.controller.queue import endpoint
2221
from octavia.i18n import _LI
2322

2423
LOG = logging.getLogger(__name__)
2524

2625

27-
class Consumer(service.Service):
26+
class ConsumerService(cotyledon.Service):
2827

29-
def __init__(self):
30-
super(Consumer, self).__init__()
31-
self.server = None
28+
def __init__(self, worker_id, conf):
29+
super(ConsumerService, self).__init__(worker_id)
30+
self.conf = conf
31+
self.topic = conf.oslo_messaging.topic
32+
self.server = conf.host
33+
self.endpoints = [endpoint.Endpoint()]
34+
self.access_policy = dispatcher.DefaultRPCAccessPolicy
35+
self.message_listener = None
3236

33-
def start(self):
34-
topic = cfg.CONF.oslo_messaging.topic
35-
server = cfg.CONF.host
36-
transport = messaging.get_transport(cfg.CONF)
37-
target = messaging.Target(topic=topic, server=server, fanout=False)
38-
endpoints = [endpoint.Endpoint()]
39-
access_policy = dispatcher.DefaultRPCAccessPolicy
40-
self.server = messaging.get_rpc_server(transport, target, endpoints,
41-
executor='eventlet',
42-
access_policy=access_policy)
37+
def run(self):
4338
LOG.info(_LI('Starting consumer...'))
44-
self.server.start()
45-
super(Consumer, self).start()
46-
47-
def stop(self, graceful=False):
48-
if self.server:
39+
transport = messaging.get_transport(self.conf)
40+
target = messaging.Target(topic=self.topic, server=self.server,
41+
fanout=False)
42+
self.message_listener = messaging.get_rpc_server(
43+
transport, target, self.endpoints,
44+
executor='threading', access_policy=self.access_policy)
45+
self.message_listener.start()
46+
47+
def terminate(self, graceful=False):
48+
if self.message_listener:
4949
LOG.info(_LI('Stopping consumer...'))
50-
self.server.stop()
50+
self.message_listener.stop()
5151
if graceful:
5252
LOG.info(
5353
_LI('Consumer successfully stopped. Waiting for final '
5454
'messages to be processed...'))
55-
self.server.wait()
56-
super(Consumer, self).stop(graceful=graceful)
57-
58-
def reset(self):
59-
if self.server:
60-
self.server.reset()
61-
super(Consumer, self).reset()
55+
self.message_listener.wait()
56+
super(ConsumerService, self).terminate()

octavia/tests/unit/controller/queue/test_consumer.py

Lines changed: 15 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -34,9 +34,10 @@ def setUp(self):
3434
conf = self.useFixture(oslo_fixture.Config(cfg.CONF))
3535
conf.config(group="oslo_messaging", topic='foo_topic')
3636
conf.config(host='test-hostname')
37+
self.conf = conf.conf
3738

38-
def test_consumer_start(self, mock_rpc_server, mock_endpoint, mock_target,
39-
mock_get_transport):
39+
def test_consumer_run(self, mock_rpc_server, mock_endpoint, mock_target,
40+
mock_get_transport):
4041
mock_get_transport_rv = mock.Mock()
4142
mock_get_transport.return_value = mock_get_transport_rv
4243
mock_rpc_server_rv = mock.Mock()
@@ -46,7 +47,7 @@ def test_consumer_start(self, mock_rpc_server, mock_endpoint, mock_target,
4647
mock_target_rv = mock.Mock()
4748
mock_target.return_value = mock_target_rv
4849

49-
consumer.Consumer().start()
50+
consumer.ConsumerService(1, self.conf).run()
5051

5152
mock_get_transport.assert_called_once_with(cfg.CONF)
5253
mock_target.assert_called_once_with(topic='foo_topic',
@@ -57,27 +58,27 @@ def test_consumer_start(self, mock_rpc_server, mock_endpoint, mock_target,
5758
mock_rpc_server.assert_called_once_with(mock_get_transport_rv,
5859
mock_target_rv,
5960
[mock_endpoint_rv],
60-
executor='eventlet',
61+
executor='threading',
6162
access_policy=access_policy)
6263

63-
def test_consumer_stop(self, mock_rpc_server, mock_endpoint, mock_target,
64-
mock_get_transport):
64+
def test_consumer_terminate(self, mock_rpc_server, mock_endpoint,
65+
mock_target, mock_get_transport):
6566
mock_rpc_server_rv = mock.Mock()
6667
mock_rpc_server.return_value = mock_rpc_server_rv
6768

68-
cons = consumer.Consumer()
69-
cons.start()
70-
cons.stop()
69+
cons = consumer.ConsumerService(1, self.conf)
70+
cons.run()
71+
cons.terminate()
7172
mock_rpc_server_rv.stop.assert_called_once_with()
7273
self.assertFalse(mock_rpc_server_rv.wait.called)
7374

74-
def test_consumer_graceful_stop(self, mock_rpc_server, mock_endpoint,
75-
mock_target, mock_get_transport):
75+
def test_consumer_graceful_terminate(self, mock_rpc_server, mock_endpoint,
76+
mock_target, mock_get_transport):
7677
mock_rpc_server_rv = mock.Mock()
7778
mock_rpc_server.return_value = mock_rpc_server_rv
7879

79-
cons = consumer.Consumer()
80-
cons.start()
81-
cons.stop(graceful=True)
80+
cons = consumer.ConsumerService(1, self.conf)
81+
cons.run()
82+
cons.terminate(graceful=True)
8283
mock_rpc_server_rv.stop.assert_called_once_with()
8384
mock_rpc_server_rv.wait.assert_called_once_with()

requirements.txt

Lines changed: 1 addition & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -2,11 +2,11 @@
22
# of appearance. Changing the order has an impact on the overall integration
33
# process, which may cause wedges in the gate later.
44
alembic>=0.8.10 # MIT
5+
cotyledon>=1.3.0 # Apache-2.0
56
pecan!=1.0.2,!=1.0.3,!=1.0.4,!=1.2,>=1.0.0 # BSD
67
pbr!=2.1.0,>=2.0.0 # Apache-2.0
78
SQLAlchemy!=1.1.5,!=1.1.6,!=1.1.7,!=1.1.8,>=1.0.10 # MIT
89
Babel!=2.4.0,>=2.3.4 # BSD
9-
eventlet!=0.18.3,>=0.18.2 # MIT
1010
requests!=2.12.2,!=2.13.0,>=2.10.0 # Apache-2.0
1111
rfc3986>=0.3.1 # Apache-2.0
1212
keystoneauth1>=2.18.0 # Apache-2.0
@@ -24,7 +24,6 @@ oslo.messaging>=5.19.0 # Apache-2.0
2424
oslo.middleware>=3.10.0 # Apache-2.0
2525
oslo.policy>=1.17.0 # Apache-2.0
2626
oslo.reports>=0.6.0 # Apache-2.0
27-
oslo.service>=1.10.0 # Apache-2.0
2827
oslo.utils>=3.20.0 # Apache-2.0
2928
pyasn1 # BSD
3029
pyasn1-modules # BSD

0 commit comments

Comments
 (0)