diff --git a/python/bullmq/flow_producer.py b/python/bullmq/flow_producer.py index 928e3f26a4f..260b4afdf12 100644 --- a/python/bullmq/flow_producer.py +++ b/python/bullmq/flow_producer.py @@ -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 @@ -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 @@ -32,10 +33,13 @@ 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 @@ -43,17 +47,17 @@ def __init__(self, redisOpts: Union[dict, str] = {}, opts: QueueBaseOptions = {} 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) @@ -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") @@ -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): """