Skip to content
Open
Changes from 1 commit
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
14 changes: 7 additions & 7 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 Down Expand Up @@ -43,17 +43,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)
Expand All @@ -62,7 +62,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,7 +116,7 @@ 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 = {}) -> dict:
parent_opts = flow.get("opts", {}).get("parent", None)

Copilot AI Apr 21, 2026

Copy link

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

opts: dict = {} uses a mutable default argument. Even if currently treated as read-only, this can lead to shared state across calls and is discouraged. Prefer opts: Optional[dict] = None (or similar) and initialize an empty dict inside the method.

Copilot uses AI. Check for mistakes.

result = None
Expand All @@ -127,7 +127,7 @@ async def add(self, flow: dict, opts: dict = {}):

return result

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