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 diff --git a/python/fluxqueue/_task.py b/python/fluxqueue/_task.py index 7f9d9b0..ee95898 100644 --- a/python/fluxqueue/_task.py +++ b/python/fluxqueue/_task.py @@ -1,11 +1,15 @@ +from __future__ import annotations + 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") @@ -47,7 +51,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 diff --git a/python/fluxqueue/context.py b/python/fluxqueue/context.py index d8686e6..fc605a9 100644 --- a/python/fluxqueue/context.py +++ b/python/fluxqueue/context.py @@ -1,12 +1,25 @@ +from __future__ import annotations + import inspect 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")