Skip to content
Open
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
33 changes: 16 additions & 17 deletions python/bullmq/flow_producer.py
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
from typing import Union
from typing import Any, Optional, Union
from bullmq.redis_connection import RedisConnection
from bullmq.types import QueueBaseOptions
from bullmq.scripts import Scripts
Expand All @@ -12,10 +12,11 @@ class MinimalQueue:
Instantiate a MinimalQueue object
"""

def __init__(self, name: str, queue_keys, redisConnection, scripts, opts: QueueBaseOptions = {}):
def __init__(self, name: str, queue_keys, redisConnection, scripts, opts: QueueBaseOptions | None = None):
"""
Initialize a connection
"""
opts = opts or {}
self.name = name
self.redisConnection = redisConnection
self.client = self.redisConnection.conn
Expand All @@ -32,28 +33,31 @@ class FlowProducer:
"""

#TODO: pass only queueOpts, no need 2 parameters in next breaking change
def __init__(self, redisOpts: Union[dict, str] = {}, opts: QueueBaseOptions = {}):
def __init__(self, redisOpts: Union[dict, str] | None = None, opts: QueueBaseOptions | None = None):
"""
Initialize a connection
"""
if redisOpts is None:
redisOpts = {}
opts = opts or {}
self.redisConnection = RedisConnection(redisOpts)
self.client = self.redisConnection.conn
self.opts: dict = opts
self.prefix = opts.get("prefix", "bull")
self.scripts = Scripts(
self.prefix, "__default__", self.redisConnection)

def queueFromNode(self, node:dict, queue_keys, prefix: str):
def queueFromNode(self, node: dict, queue_keys: QueueKeys, prefix: str) -> MinimalQueue:
return MinimalQueue(node.get("queueName"), queue_keys, self.redisConnection, self.scripts, {"prefix": prefix})

async def addChildren(self, nodes, parent, queues_opts, pipe):
async def addChildren(self, nodes: list[dict], parent: dict, queues_opts: Optional[dict], pipe: Any) -> list[dict]:
children = []
for node in nodes:
job = await self.addNode(node, parent, queues_opts, pipe)
children.append(job)
return children

async def addNodes(self, nodes: list[dict], pipe):
async def addNodes(self, nodes: list[dict], pipe: Any) -> list[dict]:
trees = []
for node in nodes:
parent_opts = node.get("opts", {}).get("parent", None)
Expand All @@ -62,7 +66,7 @@ async def addNodes(self, nodes: list[dict], pipe):

return trees

async def addNode(self, node: dict, parent: dict, queues_opts: dict, pipe):
async def addNode(self, node: dict, parent: dict, queues_opts: Optional[dict], pipe: Any) -> dict:
prefix = node.get("prefix", self.prefix)
queue = self.queueFromNode(node, QueueKeys(prefix), prefix)
queue_name = node.get("queueName")
Expand Down Expand Up @@ -116,25 +120,20 @@ async def addNode(self, node: dict, parent: dict, queues_opts: dict, pipe):

return {"job": job}

async def add(self, flow: dict, opts: dict = {}):
async def add(self, flow: dict, opts: dict | None = None) -> dict:
opts = opts or {}
parent_opts = flow.get("opts", {}).get("parent", None)

result = None
async with self.redisConnection.conn.pipeline(transaction=True) as pipe:
jobs_tree = await self.addNode(flow, {"parentOpts": parent_opts},opts.get("queuesOptions"), pipe)
await pipe.execute()
result = jobs_tree
return jobs_tree

return result

async def addBulk(self, flows: list[dict]):
result = None
async def addBulk(self, flows: list[dict]) -> list[dict]:
async with self.redisConnection.conn.pipeline(transaction=True) as pipe:
job_trees = await self.addNodes(flows, pipe)
await pipe.execute()
result = job_trees

return result
return job_trees

async def close(self):
"""
Expand Down