Merge pull request #414 from crea-asia/cluster_id

Differentiate between PID and unique cluster ID
This commit is contained in:
Ilan Steemers
2020-02-14 14:34:21 +01:00
committed by GitHub
5 changed files with 33 additions and 21 deletions
+10 -7
View File
@@ -14,6 +14,7 @@ import importlib
import signal import signal
import socket import socket
import traceback import traceback
import uuid
# Django # Django
from django import db from django import db
from django.utils import timezone from django.utils import timezone
@@ -38,6 +39,7 @@ class Cluster(object):
self.stop_event = None self.stop_event = None
self.start_event = None self.start_event = None
self.pid = current_process().pid self.pid = current_process().pid
self.cluster_id = uuid.uuid4()
self.host = socket.gethostname() self.host = socket.gethostname()
self.timeout = Conf.TIMEOUT self.timeout = Conf.TIMEOUT
signal.signal(signal.SIGTERM, self.sig_handler) signal.signal(signal.SIGTERM, self.sig_handler)
@@ -48,9 +50,9 @@ class Cluster(object):
self.stop_event = Event() self.stop_event = Event()
self.start_event = Event() self.start_event = Event()
self.sentinel = Process(target=Sentinel, self.sentinel = Process(target=Sentinel,
args=(self.stop_event, self.start_event, self.broker, self.timeout)) args=(self.stop_event, self.start_event, self.cluster_id, self.broker, self.timeout))
self.sentinel.start() self.sentinel.start()
logger.info(_('Q Cluster-{} starting.').format(self.pid)) logger.info(_('Q Cluster-{} starting.').format(self.cluster_id))
while not self.start_event.is_set(): while not self.start_event.is_set():
sleep(0.1) sleep(0.1)
return self.pid return self.pid
@@ -58,10 +60,10 @@ class Cluster(object):
def stop(self): def stop(self):
if not self.sentinel.is_alive(): if not self.sentinel.is_alive():
return False return False
logger.info(_('Q Cluster-{} stopping.').format(self.pid)) logger.info(_('Q Cluster-{} stopping.').format(self.cluster_id))
self.stop_event.set() self.stop_event.set()
self.sentinel.join() self.sentinel.join()
logger.info(_('Q Cluster-{} has stopped.').format(self.pid)) logger.info(_('Q Cluster-{} has stopped.').format(self.cluster_id))
self.start_event = None self.start_event = None
self.stop_event = None self.stop_event = None
return True return True
@@ -74,8 +76,8 @@ class Cluster(object):
@property @property
def stat(self): def stat(self):
if self.sentinel: if self.sentinel:
return Stat.get(self.pid) return Stat.get(pid=self.pid, cluster_id=self.cluster_id)
return Status(self.pid) return Status(pid=self.pid, cluster_id=self.cluster_id)
@property @property
def is_starting(self): def is_starting(self):
@@ -95,11 +97,12 @@ class Cluster(object):
class Sentinel(object): class Sentinel(object):
def __init__(self, stop_event, start_event, broker=None, timeout=Conf.TIMEOUT, start=True): def __init__(self, stop_event, start_event, cluster_id, broker=None, timeout=Conf.TIMEOUT, start=True):
# Make sure we catch signals for the pool # Make sure we catch signals for the pool
signal.signal(signal.SIGINT, signal.SIG_IGN) signal.signal(signal.SIGINT, signal.SIG_IGN)
signal.signal(signal.SIGTERM, signal.SIG_DFL) signal.signal(signal.SIGTERM, signal.SIG_DFL)
self.pid = current_process().pid self.pid = current_process().pid
self.cluster_id = cluster_id
self.parent_pid = get_ppid() self.parent_pid = get_ppid()
self.name = current_process().name self.name = current_process().name
self.broker = broker or get_broker() self.broker = broker or get_broker()
+1 -1
View File
@@ -72,7 +72,7 @@ def monitor(run_once=False, broker=None):
uptime = '%d:%02d:%02d' % (hours, minutes, seconds) uptime = '%d:%02d:%02d' % (hours, minutes, seconds)
# print to the terminal # print to the terminal
print(term.move(i, 0) + term.center(stat.host[:col_width - 1], width=col_width - 1)) print(term.move(i, 0) + term.center(stat.host[:col_width - 1], width=col_width - 1))
print(term.move(i, 1 * col_width) + term.center(stat.cluster_id, width=col_width - 1)) print(term.move(i, 1 * col_width) + term.center(str(stat.cluster_id)[-8:], width=col_width - 1))
print(term.move(i, 2 * col_width) + term.center(status, width=col_width - 1)) print(term.move(i, 2 * col_width) + term.center(status, width=col_width - 1))
print(term.move(i, 3 * col_width) + term.center(workers, width=col_width - 1)) print(term.move(i, 3 * col_width) + term.center(workers, width=col_width - 1))
print(term.move(i, 4 * col_width) + term.center(tasks, width=col_width - 1)) print(term.move(i, 4 * col_width) + term.center(tasks, width=col_width - 1))
+6 -5
View File
@@ -8,11 +8,12 @@ from django_q.signing import SignedPackage, BadSignature
class Status(object): class Status(object):
"""Cluster status base class.""" """Cluster status base class."""
def __init__(self, pid): def __init__(self, pid, cluster_id):
self.workers = [] self.workers = []
self.tob = None self.tob = None
self.reincarnations = 0 self.reincarnations = 0
self.cluster_id = pid self.pid = pid
self.cluster_id = cluster_id
self.sentinel = 0 self.sentinel = 0
self.status = Conf.STOPPED self.status = Conf.STOPPED
self.done_q_size = 0 self.done_q_size = 0
@@ -27,7 +28,7 @@ class Stat(Status):
"""Status object for Cluster monitoring.""" """Status object for Cluster monitoring."""
def __init__(self, sentinel): def __init__(self, sentinel):
super(Stat, self).__init__(sentinel.parent_pid or sentinel.pid) super(Stat, self).__init__(sentinel.parent_pid or sentinel.pid, cluster_id=sentinel.cluster_id)
self.broker = sentinel.broker or get_broker() self.broker = sentinel.broker or get_broker()
self.tob = sentinel.tob self.tob = sentinel.tob
self.reincarnations = sentinel.reincarnations self.reincarnations = sentinel.reincarnations
@@ -72,7 +73,7 @@ class Stat(Status):
return self.done_q_size + self.task_q_size == 0 return self.done_q_size + self.task_q_size == 0
@staticmethod @staticmethod
def get(cluster_id, broker=None): def get(pid, cluster_id, broker=None):
""" """
gets the current status for the cluster gets the current status for the cluster
:param cluster_id: id of the cluster :param cluster_id: id of the cluster
@@ -86,7 +87,7 @@ class Stat(Status):
return SignedPackage.loads(pack) return SignedPackage.loads(pack)
except BadSignature: except BadSignature:
return None return None
return Status(cluster_id) return Status(pid=pid, cluster_id=cluster_id)
@staticmethod @staticmethod
def get_all(broker=None): def get_all(broker=None):
+12 -6
View File
@@ -3,6 +3,7 @@ import threading
from multiprocessing import Event, Value from multiprocessing import Event, Value
from time import sleep from time import sleep
from django.utils import timezone from django.utils import timezone
import uuid as uuidlib
import os import os
import pytest import pytest
@@ -72,7 +73,8 @@ def test_sentinel():
start_event = Event() start_event = Event()
stop_event = Event() stop_event = Event()
stop_event.set() stop_event.set()
s = Sentinel(stop_event, start_event, broker=get_broker('sentinel_test:q')) cluster_id = uuidlib.uuid4()
s = Sentinel(stop_event, start_event, cluster_id=cluster_id, broker=get_broker('sentinel_test:q'))
assert start_event.is_set() assert start_event.is_set()
assert s.status() == Conf.STOPPED assert s.status() == Conf.STOPPED
@@ -256,9 +258,10 @@ def test_timeout(broker, cluster_config_timeout, async_task_kwargs):
async_task('time.sleep', 5, broker=broker, **async_task_kwargs) async_task('time.sleep', 5, broker=broker, **async_task_kwargs)
start_event = Event() start_event = Event()
stop_event = Event() stop_event = Event()
cluster_id = uuidlib.uuid4()
# Set a timer to stop the Sentinel # Set a timer to stop the Sentinel
threading.Timer(3, stop_event.set).start() threading.Timer(3, stop_event.set).start()
s = Sentinel(stop_event, start_event, broker=broker, timeout=cluster_config_timeout) s = Sentinel(stop_event, start_event, cluster_id=cluster_id, broker=broker, timeout=cluster_config_timeout)
assert start_event.is_set() assert start_event.is_set()
assert s.status() == Conf.STOPPED assert s.status() == Conf.STOPPED
assert s.reincarnations == 1 assert s.reincarnations == 1
@@ -279,9 +282,10 @@ def test_timeout_task_finishes(broker, cluster_config_timeout, async_task_kwargs
async_task('time.sleep', 3, broker=broker, **async_task_kwargs) async_task('time.sleep', 3, broker=broker, **async_task_kwargs)
start_event = Event() start_event = Event()
stop_event = Event() stop_event = Event()
cluster_id = uuidlib.uuid4()
# Set a timer to stop the Sentinel # Set a timer to stop the Sentinel
threading.Timer(6, stop_event.set).start() threading.Timer(6, stop_event.set).start()
s = Sentinel(stop_event, start_event, broker=broker, timeout=cluster_config_timeout) s = Sentinel(stop_event, start_event, cluster_id=cluster_id, broker=broker, timeout=cluster_config_timeout)
assert start_event.is_set() assert start_event.is_set()
assert s.status() == Conf.STOPPED assert s.status() == Conf.STOPPED
assert s.reincarnations == 0 assert s.reincarnations == 0
@@ -297,12 +301,13 @@ def test_recycle(broker, monkeypatch):
async_task('django_q.tests.tasks.multiply', 2, 2, broker=broker) async_task('django_q.tests.tasks.multiply', 2, 2, broker=broker)
start_event = Event() start_event = Event()
stop_event = Event() stop_event = Event()
cluster_id = uuidlib.uuid4()
# override settings # override settings
monkeypatch.setattr(Conf, 'RECYCLE', 2) monkeypatch.setattr(Conf, 'RECYCLE', 2)
monkeypatch.setattr(Conf, 'WORKERS', 1) monkeypatch.setattr(Conf, 'WORKERS', 1)
# set a timer to stop the Sentinel # set a timer to stop the Sentinel
threading.Timer(3, stop_event.set).start() threading.Timer(3, stop_event.set).start()
s = Sentinel(stop_event, start_event, broker=broker) s = Sentinel(stop_event, start_event, cluster_id=cluster_id, broker=broker)
assert start_event.is_set() assert start_event.is_set()
assert s.status() == Conf.STOPPED assert s.status() == Conf.STOPPED
assert s.reincarnations == 1 assert s.reincarnations == 1
@@ -333,13 +338,14 @@ def test_bad_secret(broker, monkeypatch):
stop_event = Event() stop_event = Event()
stop_event.set() stop_event.set()
start_event = Event() start_event = Event()
s = Sentinel(stop_event, start_event, broker=broker, start=False) cluster_id = uuidlib.uuid4()
s = Sentinel(stop_event, start_event, cluster_id=cluster_id, broker=broker, start=False)
Stat(s).save() Stat(s).save()
# change the SECRET # change the SECRET
monkeypatch.setattr(Conf, "SECRET_KEY", "OOPS") monkeypatch.setattr(Conf, "SECRET_KEY", "OOPS")
stat = Stat.get_all() stat = Stat.get_all()
assert len(stat) == 0 assert len(stat) == 0
assert Stat.get(s.parent_pid) is None assert Stat.get(pid=s.parent_pid, cluster_id=cluster_id) is None
task_queue = Queue() task_queue = Queue()
pusher(task_queue, stop_event, broker=broker) pusher(task_queue, stop_event, broker=broker)
result_queue = Queue() result_queue = Queue()
+4 -2
View File
@@ -1,4 +1,5 @@
import pytest import pytest
import uuid
from django_q.tasks import async_task from django_q.tasks import async_task
from django_q.brokers import get_broker from django_q.brokers import get_broker
@@ -10,7 +11,8 @@ from django_q.conf import Conf
@pytest.mark.django_db @pytest.mark.django_db
def test_monitor(monkeypatch): def test_monitor(monkeypatch):
assert Stat.get(0).sentinel == 0 cluster_id = uuid.uuid4()
assert Stat.get(pid=0, cluster_id=cluster_id).sentinel == 0
c = Cluster() c = Cluster()
c.start() c.start()
stats = monitor(run_once=True) stats = monitor(run_once=True)
@@ -18,7 +20,7 @@ def test_monitor(monkeypatch):
assert len(stats) > 0 assert len(stats) > 0
found_c = False found_c = False
for stat in stats: for stat in stats:
if stat.cluster_id == c.pid: if stat.cluster_id == c.cluster_id:
found_c = True found_c = True
assert stat.uptime() > 0 assert stat.uptime() > 0
assert stat.empty_queues() is True assert stat.empty_queues() is True