mirror of
https://github.com/django-q2/django-q2.git
synced 2026-09-29 18:18:12 +08:00
Rework - configuration now in class - better stop handling under heavy load. - scheduler test - better monitoring - better status reporting
This commit is contained in:
@@ -4,12 +4,14 @@ from multiprocessing import Queue, Event
|
||||
|
||||
import pytest
|
||||
|
||||
from conf import Conf
|
||||
|
||||
myPath = os.path.dirname(os.path.abspath(__file__))
|
||||
sys.path.insert(0, myPath + '/../')
|
||||
|
||||
from django_q.core import Cluster, async, pusher, worker, monitor, redis_client, Sentinel
|
||||
from django_q.core import Cluster, async, pusher, worker, monitor, redis_client, Sentinel, scheduler
|
||||
from django_q.humanhash import DEFAULT_WORDLIST
|
||||
from django_q import result, get_task, Task
|
||||
from django_q import result, get_task, Task, Schedule
|
||||
from django_q.tests.tasks import multiply
|
||||
|
||||
|
||||
@@ -30,6 +32,7 @@ def test_redis_connection(r):
|
||||
assert r.ping() is True
|
||||
|
||||
|
||||
@pytest.mark.django_db
|
||||
def test_cluster_initial():
|
||||
c = Cluster()
|
||||
assert c.sentinel is None
|
||||
@@ -37,17 +40,21 @@ def test_cluster_initial():
|
||||
assert c.start() > 0
|
||||
assert c.sentinel.is_alive() is True
|
||||
assert c.is_running
|
||||
stat = c.stat
|
||||
assert stat.status == Conf.IDLE
|
||||
assert c.stop() is True
|
||||
assert c.sentinel.is_alive() is False
|
||||
assert c.has_stopped
|
||||
|
||||
|
||||
@pytest.mark.django_db
|
||||
def test_sentinel():
|
||||
start_event = Event()
|
||||
stop_event = Event()
|
||||
stop_event.set()
|
||||
Sentinel(stop_event, start_event, list_key='sentinel_test:q')
|
||||
s = Sentinel(stop_event, start_event, list_key='sentinel_test:q')
|
||||
assert start_event.is_set()
|
||||
assert s.status() == Conf.STOPPED
|
||||
|
||||
|
||||
@pytest.mark.django_db
|
||||
@@ -165,7 +172,6 @@ def test_async(r):
|
||||
r.delete(list_key)
|
||||
|
||||
|
||||
# not sure if this actually asserts, but it is called
|
||||
@pytest.mark.django_db
|
||||
def assert_result(task):
|
||||
assert task is not None
|
||||
|
||||
@@ -0,0 +1,46 @@
|
||||
from multiprocessing import Queue, Event
|
||||
import pytest
|
||||
from django_q.core import scheduler, pusher, worker, monitor, redis_client
|
||||
from django_q import Schedule, get_task
|
||||
|
||||
@pytest.fixture
|
||||
def r():
|
||||
return redis_client
|
||||
|
||||
@pytest.mark.django_db
|
||||
def test_scheduler(r):
|
||||
list_key = 'scheduler_test:q'
|
||||
r.delete(list_key)
|
||||
schedule = Schedule.objects.create(func='math.copysign',
|
||||
args='1, -1',
|
||||
schedule_type=Schedule.ONCE,
|
||||
repeats=1,
|
||||
hook='django_q.tests.tasks.result'
|
||||
)
|
||||
assert schedule.last_run() is None
|
||||
# run scheduler
|
||||
scheduler(list_key=list_key)
|
||||
# set up the workflow
|
||||
task_queue = Queue()
|
||||
stop_event = Event()
|
||||
stop_event.set()
|
||||
# push it
|
||||
pusher(task_queue, stop_event, list_key=list_key, r=r)
|
||||
assert task_queue.qsize() == 1
|
||||
assert r.llen(list_key) == 0
|
||||
task_queue.put('STOP')
|
||||
# let a worker handle them
|
||||
result_queue = Queue()
|
||||
worker(task_queue, result_queue)
|
||||
assert result_queue.qsize() == 1
|
||||
result_queue.put('STOP')
|
||||
# store the results
|
||||
monitor(result_queue)
|
||||
assert result_queue.qsize() == 0
|
||||
schedule.refresh_from_db()
|
||||
assert schedule.repeats == 0
|
||||
assert schedule.success() is True
|
||||
task = get_task(schedule.task)
|
||||
assert task is not None
|
||||
assert task.success is True
|
||||
assert task.result < 0
|
||||
Reference in New Issue
Block a user