From af4d49ba5e2bb7f7d7e1a6af0262f102e0d6dab1 Mon Sep 17 00:00:00 2001 From: Giorgi Merebashvili Date: Fri, 27 Feb 2026 18:38:10 +0400 Subject: [PATCH 1/4] fix: Fix tasks with the base Context class --- crates/fluxqueue-worker/src/task.rs | 18 +++++++++++++++--- python/fluxqueue/context.py | 4 ++-- 2 files changed, 17 insertions(+), 5 deletions(-) diff --git a/crates/fluxqueue-worker/src/task.rs b/crates/fluxqueue-worker/src/task.rs index 0907f42..9ad458a 100644 --- a/crates/fluxqueue-worker/src/task.rs +++ b/crates/fluxqueue-worker/src/task.rs @@ -381,7 +381,7 @@ fn get_registry(module_path: &str, queue_name: &str) -> Result ) .collect::, _>>()?; - let contexts: HashMap, Arc>> = registry + let mut contexts: HashMap, Arc>> = registry .get_item("contexts")? .expect("contexts missing") .cast::() @@ -389,11 +389,14 @@ fn get_registry(module_path: &str, queue_name: &str) -> Result .iter() .filter_map(|(key, value): (Bound, Bound)| { let name: String = key.extract().ok()?; - let func: Py = value.unbind(); - Some((Arc::new(name), Arc::new(func))) + let class: Py = value.unbind(); + Some((Arc::new(name), Arc::new(class))) }) .collect(); + let (base_name, base_class) = get_base_context_class(py)?; + contexts.insert(base_name, base_class); + Ok((tasks, contexts)) })?; @@ -416,6 +419,15 @@ fn get_task_metadata(py: Python<'_>, task: Arc) -> Result> { Ok(task_metadata) } +fn get_base_context_class(py: Python<'_>) -> Result<(Arc, Arc>)> { + let module = py.import("fluxqueue.context")?; + let context_class = module.getattr("Context")?; + Ok(( + Arc::new("_Context".to_string()), + Arc::new(context_class.unbind()), + )) +} + fn is_coroutine(py: Python<'_>, func: Arc>) -> Result { let inspect = py.import("inspect")?; let is_coro: bool = inspect diff --git a/python/fluxqueue/context.py b/python/fluxqueue/context.py index 6296ef7..b402b35 100644 --- a/python/fluxqueue/context.py +++ b/python/fluxqueue/context.py @@ -29,7 +29,7 @@ class Context: to provide domain-specific resources. """ - __fluxqueue_context__: str | None = None + __fluxqueue_context__ = "_Context" def __init__(self) -> None: self._thread_local = threading.local() @@ -73,7 +73,7 @@ async def _run_async_task( self._metadata_var.reset(token) def __init_subclass__(cls) -> None: - if not cls.__fluxqueue_context__: + if not cls.__fluxqueue_context__ or cls.__fluxqueue_context__ == "_Context": cls.__fluxqueue_context__ = cls.__name__ From b75e226340197efa510e22d03928f35cac5717c8 Mon Sep 17 00:00:00 2001 From: Giorgi Merebashvili Date: Fri, 27 Feb 2026 18:42:45 +0400 Subject: [PATCH 2/4] Add typehint --- python/fluxqueue/context.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/python/fluxqueue/context.py b/python/fluxqueue/context.py index b402b35..e7a72c9 100644 --- a/python/fluxqueue/context.py +++ b/python/fluxqueue/context.py @@ -29,7 +29,7 @@ class Context: to provide domain-specific resources. """ - __fluxqueue_context__ = "_Context" + __fluxqueue_context__: str | None = "_Context" def __init__(self) -> None: self._thread_local = threading.local() From e86f299dc54d1955d0c4c727b412ff332571e4c5 Mon Sep 17 00:00:00 2001 From: Giorgi Merebashvili Date: Fri, 27 Feb 2026 18:50:19 +0400 Subject: [PATCH 3/4] Add handling edge case of a subclass named _Context --- python/fluxqueue/context.py | 3 +++ 1 file changed, 3 insertions(+) diff --git a/python/fluxqueue/context.py b/python/fluxqueue/context.py index e7a72c9..0a24742 100644 --- a/python/fluxqueue/context.py +++ b/python/fluxqueue/context.py @@ -73,6 +73,9 @@ async def _run_async_task( self._metadata_var.reset(token) def __init_subclass__(cls) -> None: + if cls.__name__ == "_Context": + raise ValueError("Subclass cannot be named '_Context'") + if not cls.__fluxqueue_context__ or cls.__fluxqueue_context__ == "_Context": cls.__fluxqueue_context__ = cls.__name__ From aeefcde7f1bf659c488d0a90f3fdc4893783d93a Mon Sep 17 00:00:00 2001 From: Giorgi Merebashvili Date: Fri, 27 Feb 2026 19:24:06 +0400 Subject: [PATCH 4/4] Better error message --- python/fluxqueue/context.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/python/fluxqueue/context.py b/python/fluxqueue/context.py index 0a24742..9ce638c 100644 --- a/python/fluxqueue/context.py +++ b/python/fluxqueue/context.py @@ -74,7 +74,7 @@ async def _run_async_task( def __init_subclass__(cls) -> None: if cls.__name__ == "_Context": - raise ValueError("Subclass cannot be named '_Context'") + raise ValueError("Context subclass cannot be named '_Context'") if not cls.__fluxqueue_context__ or cls.__fluxqueue_context__ == "_Context": cls.__fluxqueue_context__ = cls.__name__