From 31108ed0482ab5b60b1939f93581f2e9e72a6ded Mon Sep 17 00:00:00 2001 From: Giorgi Merebashvili Date: Mon, 2 Mar 2026 13:45:32 +0400 Subject: [PATCH 1/4] Add fluxqueue attrubute to task function --- python/fluxqueue/_task.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/python/fluxqueue/_task.py b/python/fluxqueue/_task.py index 7f9d9b0..64d61dd 100644 --- a/python/fluxqueue/_task.py +++ b/python/fluxqueue/_task.py @@ -47,7 +47,7 @@ def _task_decorator( task_name = get_task_name(func, name) - # TODO: Add unique identifier 'fluxqueue' just to be 100% sure + cast(Any, func).fluxqueue = True cast(Any, func).task_name = task_name cast(Any, func).queue = queue From 9fe7afa25b6aab01a0cbac285d0202779d8ba239 Mon Sep 17 00:00:00 2001 From: Giorgi Merebashvili Date: Mon, 2 Mar 2026 13:48:08 +0400 Subject: [PATCH 2/4] Update get_registry script --- crates/fluxqueue-worker/scripts/get_registry.py | 6 ++++++ 1 file changed, 6 insertions(+) diff --git a/crates/fluxqueue-worker/scripts/get_registry.py b/crates/fluxqueue-worker/scripts/get_registry.py index 51c7b3c..3db855d 100644 --- a/crates/fluxqueue-worker/scripts/get_registry.py +++ b/crates/fluxqueue-worker/scripts/get_registry.py @@ -19,8 +19,14 @@ def get_registry( # noqa: C901 registry = {"tasks": {}, "contexts": {}} for _name, obj in inspect.getmembers(module): if inspect.isfunction(obj): + is_fluxqueue_task = getattr(obj, "fluxqueue", False) + + if not is_fluxqueue_task: + continue + task_name = getattr(obj, "task_name", None) task_queue = getattr(obj, "queue", None) + if not task_queue or task_queue != queue: continue From c1a5679f82789568aba032e106cc80166f7c3c47 Mon Sep 17 00:00:00 2001 From: Giorgi Merebashvili Date: Mon, 2 Mar 2026 13:48:42 +0400 Subject: [PATCH 3/4] Reformat imports --- python/fluxqueue/_task.py | 6 ++++-- python/fluxqueue/context.py | 17 ++++++++++++++--- 2 files changed, 18 insertions(+), 5 deletions(-) diff --git a/python/fluxqueue/_task.py b/python/fluxqueue/_task.py index 64d61dd..0c1230f 100644 --- a/python/fluxqueue/_task.py +++ b/python/fluxqueue/_task.py @@ -1,11 +1,13 @@ import inspect from collections.abc import Callable, Coroutine from functools import wraps -from typing import Any, ParamSpec, cast, get_type_hints, overload +from typing import TYPE_CHECKING, Any, ParamSpec, cast, get_type_hints, overload -from ._core import FluxQueueCore from .utils import get_task_name +if TYPE_CHECKING: + from ._core import FluxQueueCore + P = ParamSpec("P") diff --git a/python/fluxqueue/context.py b/python/fluxqueue/context.py index d8686e6..8bfeb44 100644 --- a/python/fluxqueue/context.py +++ b/python/fluxqueue/context.py @@ -2,11 +2,22 @@ import threading from collections.abc import Callable, Coroutine from contextvars import ContextVar -from typing import Any, Concatenate, ParamSpec, TypeVar, cast, get_type_hints, overload +from typing import ( + TYPE_CHECKING, + Any, + Concatenate, + ParamSpec, + TypeVar, + cast, + get_type_hints, + overload, +) -from ._core import FluxQueueCore from ._task import _task_decorator -from .models import TaskMetadata + +if TYPE_CHECKING: + from ._core import FluxQueueCore + from .models import TaskMetadata P = ParamSpec("P") From 798a885efba3cff3487ca96f00e30ea7cb639e01 Mon Sep 17 00:00:00 2001 From: Giorgi Merebashvili Date: Mon, 2 Mar 2026 13:57:36 +0400 Subject: [PATCH 4/4] Fix --- python/fluxqueue/_task.py | 2 ++ python/fluxqueue/context.py | 2 ++ 2 files changed, 4 insertions(+) diff --git a/python/fluxqueue/_task.py b/python/fluxqueue/_task.py index 0c1230f..ee95898 100644 --- a/python/fluxqueue/_task.py +++ b/python/fluxqueue/_task.py @@ -1,3 +1,5 @@ +from __future__ import annotations + import inspect from collections.abc import Callable, Coroutine from functools import wraps diff --git a/python/fluxqueue/context.py b/python/fluxqueue/context.py index 8bfeb44..fc605a9 100644 --- a/python/fluxqueue/context.py +++ b/python/fluxqueue/context.py @@ -1,3 +1,5 @@ +from __future__ import annotations + import inspect import threading from collections.abc import Callable, Coroutine