Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 2 additions & 2 deletions python/pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -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]
Expand Down
71 changes: 25 additions & 46 deletions python/tests/flow_test.py
Original file line number Diff line number Diff line change
Expand Up @@ -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')
Expand Down Expand Up @@ -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(
Expand All @@ -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}
]
}
)
Expand All @@ -83,67 +82,47 @@ 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}
]
},
{
"name": 'parent-job-2',
"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 = [
Expand Down Expand Up @@ -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(
Expand All @@ -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}
]
}
)
Expand All @@ -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, {})
Expand Down
Loading