diff --git a/python/pyproject.toml b/python/pyproject.toml index f8d524fd93a..6f255499853 100644 --- a/python/pyproject.toml +++ b/python/pyproject.toml @@ -20,10 +20,10 @@ classifiers=[ ] keywords = ["python", "bullmq", "queues"] dependencies = [ - "redis ==7.4.1", + "redis ==8.1.0", "msgpack ==1.2.1", "semver ==3.0.4", - "croniter ==2.0.7" + "croniter ==6.2.4" ] [project.optional-dependencies] diff --git a/python/tests/flow_test.py b/python/tests/flow_test.py index 8dd065ba9b8..bf67c16641b 100644 --- a/python/tests/flow_test.py +++ b/python/tests/flow_test.py @@ -14,14 +14,13 @@ import unittest -queue_name = "" prefix = os.environ.get('BULLMQ_TEST_PREFIX') or "bull" class TestJob(unittest.IsolatedAsyncioTestCase): def setUp(self): print("Setting up test queue") - queueName = f"__test_queue__{uuid4().hex}" + self.queue_name = f"__test_queue__{uuid4().hex}" async def asyncTearDown(self): connection = redis.Redis(host='localhost') @@ -53,7 +52,7 @@ async def process2(job: Job, token: str): return 1 parent_worker = Worker(parent_queue_name, process2, {"prefix": prefix}) - children_worker = Worker(queue_name, process1, {"prefix": prefix}) + children_worker = Worker(self.queue_name, process1, {"prefix": prefix}) flow = FlowProducer({"prefix": prefix}) await flow.add( @@ -62,9 +61,9 @@ async def process2(job: Job, token: str): "queueName": parent_queue_name, "data": {}, "children": [ - {"name": child_job_name, "data": {"idx": 0, "foo": 'bar'}, "queueName": queue_name}, - {"name": child_job_name, "data": {"idx": 1, "foo": 'baz'}, "queueName": queue_name}, - {"name": child_job_name, "data": {"idx": 2, "foo": 'qux'}, "queueName": queue_name} + {"name": child_job_name, "data": {"idx": 0, "foo": 'bar'}, "queueName": self.queue_name}, + {"name": child_job_name, "data": {"idx": 1, "foo": 'baz'}, "queueName": self.queue_name}, + {"name": child_job_name, "data": {"idx": 2, "foo": 'qux'}, "queueName": self.queue_name} ] } ) @@ -83,43 +82,17 @@ async def process2(job: Job, token: str): async def test_addBulk_should_process_children_before_parent(self): child_job_name = 'child-job' - children_data = [ - {"idx": 0, "bar": 'something'}, - {"idx": 1, "baz": 'something'} - ] + child_queue_name = f"__test_child_queue__{uuid4().hex}" parent_queue_name = f"__test_parent_queue__{uuid4().hex}" - processing_children = Future() - - processed_children = 0 - async def process1(job: Job, token: str): - nonlocal processed_children - processed_children+=1 - if processed_children == len(children_data): - processing_children.set_result(None) - return children_data[job.data.get("idx")] - - processing_parents = Future() - - processed_parents = 0 - async def process2(job: Job, token: str): - nonlocal processed_parents - processed_parents+=1 - if processed_parents == 2: - processing_parents.set_result(None) - return 1 - - parent_worker = Worker(parent_queue_name, process2, {"prefix": prefix}) - children_worker = Worker(queue_name, process1, {"prefix": prefix}) - flow = FlowProducer({"prefix": prefix}) - await flow.addBulk([ + trees = await flow.addBulk([ { "name": 'parent-job-1', "queueName": parent_queue_name, "data": {}, "children": [ - {"name": child_job_name, "data": {"idx": 0, "foo": 'bar'}, "queueName": queue_name} + {"name": child_job_name, "data": {"idx": 0, "foo": 'bar'}, "queueName": child_queue_name} ] }, { @@ -127,23 +100,29 @@ async def process2(job: Job, token: str): "queueName": parent_queue_name, "data": {}, "children": [ - {"name": child_job_name, "data": {"idx": 1, "foo": 'baz'}, "queueName": queue_name} + {"name": child_job_name, "data": {"idx": 1, "foo": 'baz'}, "queueName": child_queue_name} ] } ]) - await processing_children - await processing_parents + parent_queue = Queue(parent_queue_name, {"prefix": prefix}) + child_queue = Queue(child_queue_name, {"prefix": prefix}) + + for tree in trees: + self.assertTrue(await tree["job"].isWaitingChildren()) + child_job = tree["children"][0]["job"] + self.assertEqual(await child_job.getState(), "waiting") - await parent_worker.close() - await children_worker.close() await flow.close() - parent_queue = Queue(parent_queue_name, {"prefix": prefix}) await parent_queue.pause() await parent_queue.obliterate() await parent_queue.close() + await child_queue.pause() + await child_queue.obliterate() + await child_queue.close() + async def test_get_children_values(self): child_job_name = 'child-job' children_data = [ @@ -171,7 +150,7 @@ async def process2(job: Job, token: str): return 1 parent_worker = Worker(parent_queue_name, process2, {"prefix": prefix}) - children_worker = Worker(queue_name, process1, {"prefix": prefix}) + children_worker = Worker(self.queue_name, process1, {"prefix": prefix}) flow = FlowProducer({"prefix": prefix}) await flow.add( @@ -180,9 +159,9 @@ async def process2(job: Job, token: str): "queueName": parent_queue_name, "data": {}, "children": [ - {"name": child_job_name, "data": {"idx": 0, "foo": 'bar'}, "queueName": queue_name}, - {"name": child_job_name, "data": {"idx": 1, "foo": 'baz'}, "queueName": queue_name}, - {"name": child_job_name, "data": {"idx": 2, "foo": 'qux'}, "queueName": queue_name} + {"name": child_job_name, "data": {"idx": 0, "foo": 'bar'}, "queueName": self.queue_name}, + {"name": child_job_name, "data": {"idx": 1, "foo": 'baz'}, "queueName": self.queue_name}, + {"name": child_job_name, "data": {"idx": 2, "foo": 'qux'}, "queueName": self.queue_name} ] } ) @@ -207,7 +186,7 @@ def on_parent_processed(future): await parent_queue.close() async def test_get_children_values_on_simple_jobs(self): - queue = Queue(queue_name, {"prefix": prefix}) + queue = Queue(self.queue_name, {"prefix": prefix}) job = await queue.add("test", {"foo": "bar"}, {"delay": 1500}) children_values = await job.getChildrenValues() self.assertEqual(children_values, {})