mirror of
https://github.com/django-q2/django-q2.git
synced 2026-10-08 11:08:12 +08:00
simple test for monitor
This commit is contained in:
+5
-2
@@ -64,6 +64,7 @@ class Cluster(object):
|
|||||||
self.start_event = None
|
self.start_event = None
|
||||||
self.stopped_event = None
|
self.stopped_event = None
|
||||||
self.pid = current_process().pid
|
self.pid = current_process().pid
|
||||||
|
self.host = socket.gethostname()
|
||||||
self.list_key = list_key
|
self.list_key = list_key
|
||||||
signal.signal(signal.SIGTERM, self.sig_handler)
|
signal.signal(signal.SIGTERM, self.sig_handler)
|
||||||
signal.signal(signal.SIGINT, self.sig_handler)
|
signal.signal(signal.SIGINT, self.sig_handler)
|
||||||
@@ -81,6 +82,8 @@ class Cluster(object):
|
|||||||
self.sentinel = Process(target=Sentinel, args=(self.stop_event, self.start_event, self.list_key))
|
self.sentinel = Process(target=Sentinel, args=(self.stop_event, self.start_event, self.list_key))
|
||||||
self.sentinel.start()
|
self.sentinel.start()
|
||||||
logger.info('Q Cluster-{} starting.'.format(self.pid))
|
logger.info('Q Cluster-{} starting.'.format(self.pid))
|
||||||
|
while not self.start_event.is_set():
|
||||||
|
sleep(0.2)
|
||||||
return self.pid
|
return self.pid
|
||||||
|
|
||||||
def stop(self):
|
def stop(self):
|
||||||
@@ -329,7 +332,7 @@ def async(func, *args, **kwargs):
|
|||||||
|
|
||||||
class SignedPackage(object):
|
class SignedPackage(object):
|
||||||
"""
|
"""
|
||||||
Wraps Django's signing module with custom JsonPickle serializer
|
Wraps Django's signing module with custom Pickle serializer
|
||||||
"""
|
"""
|
||||||
|
|
||||||
@staticmethod
|
@staticmethod
|
||||||
@@ -350,7 +353,7 @@ class SignedPackage(object):
|
|||||||
|
|
||||||
class PickleSerializer(object):
|
class PickleSerializer(object):
|
||||||
"""
|
"""
|
||||||
Simple wrapper around JsonPickle for signing.dumps and
|
Simple wrapper around Pickle for signing.dumps and
|
||||||
signing.loads.
|
signing.loads.
|
||||||
"""
|
"""
|
||||||
|
|
||||||
|
|||||||
@@ -12,51 +12,65 @@ from django_q.core import Stat, RUNNING, STOPPED
|
|||||||
class Command(BaseCommand):
|
class Command(BaseCommand):
|
||||||
help = "Monitors cluster activity"
|
help = "Monitors cluster activity"
|
||||||
|
|
||||||
|
def add_arguments(self, parser):
|
||||||
|
parser.add_argument('--run-once',
|
||||||
|
action='store_true',
|
||||||
|
dest='run_once',
|
||||||
|
default=False,
|
||||||
|
help='Run once till first stat')
|
||||||
|
|
||||||
def handle(self, *args, **options):
|
def handle(self, *args, **options):
|
||||||
term = Terminal()
|
monitor(run_once=options['run_once'])
|
||||||
with term.fullscreen(), term.hidden_cursor(), term.cbreak():
|
|
||||||
val = None
|
|
||||||
start_width = int(term.width / 8)
|
def monitor(run_once=False):
|
||||||
while val not in (u'q', u'Q',):
|
term = Terminal()
|
||||||
col_width = int(term.width / 8)
|
with term.fullscreen(), term.hidden_cursor(), term.cbreak():
|
||||||
# In case of resize
|
val = None
|
||||||
if col_width != start_width:
|
start_width = int(term.width / 8)
|
||||||
print(term.clear)
|
while val not in (u'q', u'Q',):
|
||||||
start_width = col_width
|
col_width = int(term.width / 8)
|
||||||
print(term.move(0, 0) + term.black_on_green(term.center('Host', width=col_width - 1)))
|
# In case of resize
|
||||||
print(term.move(0, 1 * col_width) + term.black_on_green(term.center('Id', width=col_width - 1)))
|
if col_width != start_width:
|
||||||
print(term.move(0, 2 * col_width) + term.black_on_green(term.center('Status', width=col_width - 1)))
|
print(term.clear)
|
||||||
print(term.move(0, 3 * col_width) + term.black_on_green(term.center('Pool', width=col_width - 1)))
|
start_width = col_width
|
||||||
print(term.move(0, 4 * col_width) + term.black_on_green(term.center('TQ', width=col_width - 1)))
|
print(term.move(0, 0) + term.black_on_green(term.center('Host', width=col_width - 1)))
|
||||||
print(term.move(0, 5 * col_width) + term.black_on_green(term.center('RQ', width=col_width - 1)))
|
print(term.move(0, 1 * col_width) + term.black_on_green(term.center('Id', width=col_width - 1)))
|
||||||
print(term.move(0, 6 * col_width) + term.black_on_green(term.center('Deaths', width=col_width - 1)))
|
print(term.move(0, 2 * col_width) + term.black_on_green(term.center('Status', width=col_width - 1)))
|
||||||
print(term.move(0, 7 * col_width) + term.black_on_green(term.center('Uptime', width=col_width - 1)))
|
print(term.move(0, 3 * col_width) + term.black_on_green(term.center('Pool', width=col_width - 1)))
|
||||||
i = 2
|
print(term.move(0, 4 * col_width) + term.black_on_green(term.center('TQ', width=col_width - 1)))
|
||||||
stats = Stat.get_all()
|
print(term.move(0, 5 * col_width) + term.black_on_green(term.center('RQ', width=col_width - 1)))
|
||||||
print(term.clear_eos())
|
print(term.move(0, 6 * col_width) + term.black_on_green(term.center('Deaths', width=col_width - 1)))
|
||||||
for stat in stats:
|
print(term.move(0, 7 * col_width) + term.black_on_green(term.center('Uptime', width=col_width - 1)))
|
||||||
# color status
|
i = 2
|
||||||
if stat.status == RUNNING:
|
stats = Stat.get_all()
|
||||||
status = term.green(RUNNING)
|
print(term.clear_eos())
|
||||||
elif stat.status == STOPPED:
|
for stat in stats:
|
||||||
status = term.red(STOPPED)
|
# color status
|
||||||
else:
|
if stat.status == RUNNING:
|
||||||
status = term.yellow(stat.status)
|
status = term.green(RUNNING)
|
||||||
# format uptime
|
elif stat.status == STOPPED:
|
||||||
uptime = (timezone.now() - stat.tob).total_seconds()
|
status = term.red(STOPPED)
|
||||||
hours, remainder = divmod(uptime, 3600)
|
else:
|
||||||
minutes, seconds = divmod(remainder, 60)
|
status = term.yellow(stat.status)
|
||||||
uptime = '%d:%02d:%02d' % (hours, minutes, seconds)
|
# format uptime
|
||||||
print(term.move(i, 0) + term.center('{}'.format(stat.host), width=col_width - 1))
|
uptime = (timezone.now() - stat.tob).total_seconds()
|
||||||
print(term.move(i, 1 * col_width) + term.center('{}'.format(stat.cluster_id), width=col_width - 1))
|
hours, remainder = divmod(uptime, 3600)
|
||||||
print(term.move(i, 2 * col_width) + term.center('{}'.format(status), width=col_width - 1))
|
minutes, seconds = divmod(remainder, 60)
|
||||||
print(
|
uptime = '%d:%02d:%02d' % (hours, minutes, seconds)
|
||||||
term.move(i, 3 * col_width) + term.center('{}'.format(len(stat.workers)), width=col_width - 1))
|
print(term.move(i, 0) + term.center('{}'.format(stat.host), width=col_width - 1))
|
||||||
print(term.move(i, 4 * col_width) + term.center('{}'.format(stat.task_q_size), width=col_width - 1))
|
print(term.move(i, 1 * col_width) + term.center('{}'.format(stat.cluster_id), width=col_width - 1))
|
||||||
print(term.move(i, 5 * col_width) + term.center('{}'.format(stat.done_q_size), width=col_width - 1))
|
print(term.move(i, 2 * col_width) + term.center('{}'.format(status), width=col_width - 1))
|
||||||
print(term.move(i, 6 * col_width) + term.center('{}'.format(stat.reincarnations),
|
print(
|
||||||
width=col_width - 1))
|
term.move(i, 3 * col_width) + term.center('{}'.format(len(stat.workers)), width=col_width - 1))
|
||||||
print(term.move(i, 7 * col_width) + term.center('{}'.format(uptime), width=col_width - 1))
|
print(term.move(i, 4 * col_width) + term.center('{}'.format(stat.task_q_size), width=col_width - 1))
|
||||||
i += 1
|
print(term.move(i, 5 * col_width) + term.center('{}'.format(stat.done_q_size), width=col_width - 1))
|
||||||
print(term.move(i + 2, 0) + term.center('[Press q to quit]'))
|
print(term.move(i, 6 * col_width) + term.center('{}'.format(stat.reincarnations),
|
||||||
val = term.inkey(timeout=1)
|
width=col_width - 1))
|
||||||
|
print(term.move(i, 7 * col_width) + term.center('{}'.format(uptime), width=col_width - 1))
|
||||||
|
i += 1
|
||||||
|
# for testing
|
||||||
|
if run_once:
|
||||||
|
return Stat.get_all()
|
||||||
|
print(term.move(i + 2, 0) + term.center('[Press q to quit]'))
|
||||||
|
val = term.inkey(timeout=1)
|
||||||
|
|||||||
@@ -0,0 +1,7 @@
|
|||||||
|
def test_admin_view(admin_client):
|
||||||
|
response = admin_client.get('/admin/django_q/')
|
||||||
|
assert response.status_code == 200
|
||||||
|
response = admin_client.get('/admin/django_q/failure/')
|
||||||
|
assert response.status_code == 200
|
||||||
|
response = admin_client.get('/admin/django_q/success/')
|
||||||
|
assert response.status_code == 200
|
||||||
@@ -20,16 +20,6 @@ class WordClass(object):
|
|||||||
def get_words(self):
|
def get_words(self):
|
||||||
return self.word_list
|
return self.word_list
|
||||||
|
|
||||||
|
|
||||||
def test_admin_view(admin_client):
|
|
||||||
response = admin_client.get('/admin/django_q/')
|
|
||||||
assert response.status_code == 200
|
|
||||||
response = admin_client.get('/admin/django_q/failure/')
|
|
||||||
assert response.status_code == 200
|
|
||||||
response = admin_client.get('/admin/django_q/success/')
|
|
||||||
assert response.status_code == 200
|
|
||||||
|
|
||||||
|
|
||||||
def test_redis_connection():
|
def test_redis_connection():
|
||||||
assert r.ping() is True
|
assert r.ping() is True
|
||||||
|
|
||||||
@@ -39,13 +29,9 @@ def test_cluster_initial():
|
|||||||
assert c.sentinel is None
|
assert c.sentinel is None
|
||||||
assert c.is_idle
|
assert c.is_idle
|
||||||
c.start()
|
c.start()
|
||||||
while c.is_starting:
|
|
||||||
sleep(0.2)
|
|
||||||
assert c.sentinel.is_alive() is True
|
assert c.sentinel.is_alive() is True
|
||||||
assert c.is_running
|
assert c.is_running
|
||||||
c.stop()
|
c.stop()
|
||||||
while c.is_stopping:
|
|
||||||
sleep(0.2)
|
|
||||||
assert c.sentinel.is_alive() is False
|
assert c.sentinel.is_alive() is False
|
||||||
assert c.has_stopped
|
assert c.has_stopped
|
||||||
|
|
||||||
@@ -92,8 +78,6 @@ def run_cluster():
|
|||||||
r.delete(list_key)
|
r.delete(list_key)
|
||||||
c = Cluster(list_key=list_key)
|
c = Cluster(list_key=list_key)
|
||||||
assert c.start() > 0
|
assert c.start() > 0
|
||||||
while not c.is_running:
|
|
||||||
sleep(0.5)
|
|
||||||
while c.stat.task_q_size > 0 and c.stat.done_q_size > 0:
|
while c.stat.task_q_size > 0 and c.stat.done_q_size > 0:
|
||||||
sleep(0.5)
|
sleep(0.5)
|
||||||
assert c.stop() is True
|
assert c.stop() is True
|
||||||
@@ -0,0 +1,11 @@
|
|||||||
|
from django_q.core import Cluster
|
||||||
|
from django_q.management.commands.qmonitor import monitor
|
||||||
|
|
||||||
|
|
||||||
|
def test_monitor():
|
||||||
|
c = Cluster()
|
||||||
|
c.start()
|
||||||
|
stats = monitor(run_once=True)
|
||||||
|
c.stop()
|
||||||
|
assert len(stats) > 0
|
||||||
|
assert stats[0].cluster_id == c.pid
|
||||||
@@ -36,6 +36,7 @@ setup(
|
|||||||
description='A multiprocessing task queue for Django',
|
description='A multiprocessing task queue for Django',
|
||||||
long_description=README,
|
long_description=README,
|
||||||
install_requires=['django>=1.7', 'redis', 'coloredlogs', 'django-picklefield', 'blessed'],
|
install_requires=['django>=1.7', 'redis', 'coloredlogs', 'django-picklefield', 'blessed'],
|
||||||
|
test_requires=['pytest', 'pytest-django', ],
|
||||||
cmdclass={'test': PyTest},
|
cmdclass={'test': PyTest},
|
||||||
classifiers=[
|
classifiers=[
|
||||||
'Development Status :: 2 - PreAlpha',
|
'Development Status :: 2 - PreAlpha',
|
||||||
|
|||||||
Reference in New Issue
Block a user