From a11abfc82931a64b979be9cfab8d268ffde5ebbf Mon Sep 17 00:00:00 2001 From: Giorgi Merebashvili Date: Sat, 28 Feb 2026 18:37:32 +0400 Subject: [PATCH 1/9] Fix async example in README --- README.md | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/README.md b/README.md index 4b59a09..fb2fd3c 100644 --- a/README.md +++ b/README.md @@ -78,7 +78,7 @@ FluxQueue supports async functions too. Just define an async function and use th ```python @fluxqueue.task() -async def send_email(data: dict): +async def send_email(to_email: str, subject: str, body: str): async with email_context() as email_client: message = EmailMessage() message["From"] = "test@example.com" From 93e9ebe451a772c8eee0555ec6ac8d6a6f74e28d Mon Sep 17 00:00:00 2001 From: Giorgi Merebashvili Date: Sat, 28 Feb 2026 18:47:05 +0400 Subject: [PATCH 2/9] Fix type checker error in __init__ --- python/fluxqueue/__init__.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/python/fluxqueue/__init__.py b/python/fluxqueue/__init__.py index de2f3ed..299142e 100644 --- a/python/fluxqueue/__init__.py +++ b/python/fluxqueue/__init__.py @@ -1,4 +1,4 @@ -from ._core import __version__ as __version__ +from ._core import __version__ as __version__ # type: ignore from .client import FluxQueue as FluxQueue from .context import Context as Context from .models import TaskMetadata as TaskMetadata From 05b82c1b24e5f505da0938ca00088468408e48d0 Mon Sep 17 00:00:00 2001 From: Giorgi Merebashvili Date: Sat, 28 Feb 2026 18:48:28 +0400 Subject: [PATCH 3/9] Remove unused T --- python/fluxqueue/_core.pyi | 4 +--- 1 file changed, 1 insertion(+), 3 deletions(-) diff --git a/python/fluxqueue/_core.pyi b/python/fluxqueue/_core.pyi index 4ec946d..a1862d3 100644 --- a/python/fluxqueue/_core.pyi +++ b/python/fluxqueue/_core.pyi @@ -1,6 +1,4 @@ -from typing import Any, TypeVar - -T = TypeVar("T") +from typing import Any class FluxQueueCore: """ From 498fb0fd46729c2f78cde8e6eb90357072590087 Mon Sep 17 00:00:00 2001 From: Giorgi Merebashvili Date: Sat, 28 Feb 2026 18:56:13 +0400 Subject: [PATCH 4/9] Add client library tests --- tests/conftest.py | 4 ++-- tests/test_context.py | 43 +++++++++++++++++++++++++++++++++++++++++++ 2 files changed, 45 insertions(+), 2 deletions(-) create mode 100644 tests/test_context.py diff --git a/tests/conftest.py b/tests/conftest.py index 43af84d..258b150 100644 --- a/tests/conftest.py +++ b/tests/conftest.py @@ -17,8 +17,8 @@ class TestEnv: @pytest.fixture def test_env(): - redis_client: Redis = Redis() - fluxqueue: FluxQueue = FluxQueue() + redis_client = Redis() + fluxqueue = FluxQueue() try: yield TestEnv(fluxqueue=fluxqueue, redis_client=redis_client) diff --git a/tests/test_context.py b/tests/test_context.py new file mode 100644 index 0000000..dca0a48 --- /dev/null +++ b/tests/test_context.py @@ -0,0 +1,43 @@ +import pytest +from fluxqueue import Context + +from .conftest import TestEnvFixture + + +def test_sync_task_with_context(test_env: TestEnvFixture): + @test_env.fluxqueue.task_with_context() + def task(ctx: Context, name: str): + print(ctx.metadata) + print("Hello ", name) + + result = task("George") + assert result is None + + test_env.redis_client.flushdb() + + +@pytest.mark.asyncio +async def test_async_task_with_context(test_env: TestEnvFixture): + @test_env.fluxqueue.task_with_context() + async def task(ctx: Context, name: str): + print(ctx.metadata) + print("Async Hello ", name) + + result = await task("Async George") + assert result is None + + test_env.redis_client.flushdb() + + +def test_task_with_context_but_without_argument(test_env: TestEnvFixture): + with pytest.raises(TypeError): + + @test_env.fluxqueue.task_with_context() # type: ignore + def task(name: str): + print("Hello ", name) + + with pytest.raises(TypeError): + + @test_env.fluxqueue.task_with_context() # type: ignore + async def async_task(name: str): + print("Async Hello ", name) From 44533c45c30cc48c710d0ed05b4bf3d2a7970f9e Mon Sep 17 00:00:00 2001 From: Giorgi Merebashvili Date: Sat, 28 Feb 2026 18:57:26 +0400 Subject: [PATCH 5/9] More detailed client tests --- tests/test_context.py | 6 ++++++ 1 file changed, 6 insertions(+) diff --git a/tests/test_context.py b/tests/test_context.py index dca0a48..b48e6f0 100644 --- a/tests/test_context.py +++ b/tests/test_context.py @@ -11,7 +11,10 @@ def task(ctx: Context, name: str): print("Hello ", name) result = task("George") + redis_result = test_env.redis_client.lrange("fluxqueue:queue:default", 0, -1) + assert result is None + assert b"George" in redis_result[0] # type: ignore test_env.redis_client.flushdb() @@ -24,7 +27,10 @@ async def task(ctx: Context, name: str): print("Async Hello ", name) result = await task("Async George") + redis_result = test_env.redis_client.lrange("fluxqueue:queue:default", 0, -1) + assert result is None + assert b"Async George" in redis_result[0] # type: ignore test_env.redis_client.flushdb() From babb5ae2b61224305f154a166a2f683274db110e Mon Sep 17 00:00:00 2001 From: Giorgi Merebashvili Date: Sat, 28 Feb 2026 19:36:07 +0400 Subject: [PATCH 6/9] Improve worker version check and add tests for it --- crates/fluxqueue-worker/src/lib.rs | 1 + crates/fluxqueue-worker/src/version_check.rs | 188 +++++++++++++++++++ crates/fluxqueue-worker/src/worker.rs | 21 +-- 3 files changed, 190 insertions(+), 20 deletions(-) create mode 100644 crates/fluxqueue-worker/src/version_check.rs diff --git a/crates/fluxqueue-worker/src/lib.rs b/crates/fluxqueue-worker/src/lib.rs index 2ae7876..a9ab96a 100644 --- a/crates/fluxqueue-worker/src/lib.rs +++ b/crates/fluxqueue-worker/src/lib.rs @@ -1,6 +1,7 @@ mod logger; mod redis_client; mod task; +mod version_check; mod worker; pub use worker::*; diff --git a/crates/fluxqueue-worker/src/version_check.rs b/crates/fluxqueue-worker/src/version_check.rs new file mode 100644 index 0000000..827803b --- /dev/null +++ b/crates/fluxqueue-worker/src/version_check.rs @@ -0,0 +1,188 @@ +fn normalize(version: &str) -> String { + version + .replace("-rc-", "rc") + .replace("-beta-", "b") + .replace("-alpha-", "a") +} + +fn parse_part(part: &str) -> (u32, i32, u32) { + let mut digits = String::new(); + let mut suffix = String::new(); + + for c in part.chars() { + if c.is_ascii_digit() && suffix.is_empty() { + digits.push(c); + } else { + suffix.push(c); + } + } + + let number = digits.parse().unwrap_or(0); + + let (ptype, pnum) = if suffix.starts_with("rc") { + (2, suffix[2..].parse().unwrap_or(0)) + } else if suffix.starts_with('b') { + (1, suffix[1..].parse().unwrap_or(0)) + } else if suffix.starts_with('a') { + (0, suffix[1..].parse().unwrap_or(0)) + } else { + (3, 0) // stable + }; + + (number, ptype, pnum) +} + +pub fn compare_versions(v1: &str, v2: &str) -> i8 { + let v1 = normalize(v1); + let v2 = normalize(v2); + + let parts1: Vec<_> = v1.split('.').collect(); + let parts2: Vec<_> = v2.split('.').collect(); + + let len = parts1.len().max(parts2.len()); + + for i in 0..len { + let p1 = parts1.get(i).unwrap_or(&"0"); + let p2 = parts2.get(i).unwrap_or(&"0"); + + let (n1, t1, r1) = parse_part(p1); + let (n2, t2, r2) = parse_part(p2); + + if n1 != n2 { + return if n1 > n2 { 1 } else { -1 }; + } + + if t1 != t2 { + return if t1 > t2 { 1 } else { -1 }; + } + + if r1 != r2 { + return if r1 > r2 { 1 } else { -1 }; + } + } + + 0 +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn test_compare_versions() { + let worker_version = "0.3.1"; + let client_lib_version = "0.1.2"; + + assert!( + compare_versions(worker_version, client_lib_version) == 1, + "left version is greater than the right one" + ); + + assert!( + compare_versions(client_lib_version, worker_version) == -1, + "left version is less than the right one" + ); + + assert!( + compare_versions("0.2.0", "0.2.0") == 0, + "both are equal versions" + ); + + assert!( + compare_versions("1.2.0", "1.2") == 0, + "both are equal versions" + ); + + assert!( + compare_versions("1.2.3", "1.2.3a1") == 1, + "alpha version is less than the latest" + ); + + assert!( + compare_versions("1.2.3", "1.2.3b1") == 1, + "beta version is less than the latest" + ); + + assert!( + compare_versions("1.2.3", "1.2.3rc1") == 1, + "rc version is less than the latest" + ); + + assert!( + compare_versions("1.2.3-rc-1", "1.2.3rc1") == 0, + "both are equal" + ); + + assert!( + compare_versions("1.2.3-beta-1", "1.2.3b1") == 0, + "both are equal" + ); + + assert!( + compare_versions("1.2.3-alpha-1", "1.2.3a1") == 0, + "both are equal" + ); + + assert!( + compare_versions("1.2.3-rc-1", "1.2.3b1") == 1, + "left is greater" + ); + + assert!(compare_versions("1.2.3a1", "1.2.3b1") == -1, "alpha < beta"); + + assert!(compare_versions("1.2.3b1", "1.2.3rc1") == -1, "beta < rc"); + + assert!(compare_versions("1.2.3a1", "1.2.3rc1") == -1, "alpha < rc"); + + assert!(compare_versions("1.2.3rc2", "1.2.3rc1") == 1, "rc2 > rc1"); + + assert!(compare_versions("1.2.3b2", "1.2.3b10") == -1, "b2 < b10"); + + assert!(compare_versions("1.2.3a10", "1.2.3a2") == 1, "a10 > a2"); + + assert!( + compare_versions("1.2.3.0", "1.2.3") == 0, + "trailing zeros ignored" + ); + + assert!( + compare_versions("1.2.3.1", "1.2.3") == 1, + "extra segment greater" + ); + + assert!( + compare_versions("1.2", "1.2.1") == -1, + "missing segment treated as zero" + ); + + assert!( + compare_versions("1.2.3-rc-2", "1.2.3rc1") == 1, + "rust rc2 > python rc1" + ); + + assert!( + compare_versions("1.2.3-beta-2", "1.2.3b1") == 1, + "rust beta2 > python b1" + ); + + assert!( + compare_versions("1.2.3-alpha-2", "1.2.3a1") == 1, + "rust alpha2 > python a1" + ); + + assert!( + compare_versions("1.2.3", "1.2.4a1") == -1, + "next version prerelease is higher" + ); + + assert!( + compare_versions("1.2.4a1", "1.2.3") == 1, + "next version prerelease is higher" + ); + + assert!( + compare_versions("1.2.3rc1", "1.2.3-rc-1") == 0, + "symmetry equality" + ); + } +} diff --git a/crates/fluxqueue-worker/src/worker.rs b/crates/fluxqueue-worker/src/worker.rs index 0bcf4e8..897373d 100644 --- a/crates/fluxqueue-worker/src/worker.rs +++ b/crates/fluxqueue-worker/src/worker.rs @@ -346,7 +346,7 @@ fn check_client_library_version() -> Result<()> { Ok(version.to_string()) })?; - let comparison = compare_versions(worker_version, &library_version); + let comparison = crate::version_check::compare_versions(worker_version, &library_version); if comparison == -1 { return Err(anyhow!( @@ -367,25 +367,6 @@ fn check_client_library_version() -> Result<()> { Ok(()) } -fn compare_versions(v1: &str, v2: &str) -> i8 { - let mut parts1: Vec = v1.split('.').map(|p| p.parse().unwrap_or(0)).collect(); - let mut parts2: Vec = v2.split('.').map(|p| p.parse().unwrap_or(0)).collect(); - - let len = parts1.len().max(parts2.len()); - parts1.resize(len, 0); - parts2.resize(len, 0); - - for (a, b) in parts1.iter().zip(parts2.iter()) { - if a > b { - return 1; - } - if a < b { - return -1; - } - } - 0 -} - #[cfg(test)] mod tests { use super::*; From 97db85971b7c363222529faa53cf080a3af6cade Mon Sep 17 00:00:00 2001 From: Giorgi Merebashvili Date: Sat, 28 Feb 2026 19:40:18 +0400 Subject: [PATCH 7/9] Fix clippy error --- crates/fluxqueue-worker/src/version_check.rs | 12 ++++++------ 1 file changed, 6 insertions(+), 6 deletions(-) diff --git a/crates/fluxqueue-worker/src/version_check.rs b/crates/fluxqueue-worker/src/version_check.rs index 827803b..a98a918 100644 --- a/crates/fluxqueue-worker/src/version_check.rs +++ b/crates/fluxqueue-worker/src/version_check.rs @@ -19,12 +19,12 @@ fn parse_part(part: &str) -> (u32, i32, u32) { let number = digits.parse().unwrap_or(0); - let (ptype, pnum) = if suffix.starts_with("rc") { - (2, suffix[2..].parse().unwrap_or(0)) - } else if suffix.starts_with('b') { - (1, suffix[1..].parse().unwrap_or(0)) - } else if suffix.starts_with('a') { - (0, suffix[1..].parse().unwrap_or(0)) + let (ptype, pnum) = if let Some(stripped) = suffix.strip_prefix("rc") { + (2, stripped.parse().unwrap_or(0)) + } else if let Some(stripped) = suffix.strip_prefix('b') { + (1, stripped.parse().unwrap_or(0)) + } else if let Some(stripped) = suffix.strip_prefix('a') { + (0, stripped.parse().unwrap_or(0)) } else { (3, 0) // stable }; From 1dc1d1c5f396bfd982bc32846d86bd327a756d66 Mon Sep 17 00:00:00 2001 From: Giorgi Merebashvili Date: Sat, 28 Feb 2026 23:31:58 +0400 Subject: [PATCH 8/9] Add more tests --- crates/fluxqueue-worker/tests/main.py | 3 +++ .../tests/test_tasks_duplicate.py | 15 +++++------ .../tests/test_tasks_module.py | 25 ++++++------------- .../tests/test_tasks_with_context.py | 8 ++++++ python/fluxqueue/__init__.py | 2 +- python/fluxqueue/_core.pyi | 2 ++ 6 files changed, 28 insertions(+), 27 deletions(-) create mode 100644 crates/fluxqueue-worker/tests/main.py create mode 100644 crates/fluxqueue-worker/tests/test_tasks_with_context.py diff --git a/crates/fluxqueue-worker/tests/main.py b/crates/fluxqueue-worker/tests/main.py new file mode 100644 index 0000000..5be352f --- /dev/null +++ b/crates/fluxqueue-worker/tests/main.py @@ -0,0 +1,3 @@ +from fluxqueue import FluxQueue + +fluxqueue = FluxQueue() diff --git a/crates/fluxqueue-worker/tests/test_tasks_duplicate.py b/crates/fluxqueue-worker/tests/test_tasks_duplicate.py index ba854fd..ba09fc7 100644 --- a/crates/fluxqueue-worker/tests/test_tasks_duplicate.py +++ b/crates/fluxqueue-worker/tests/test_tasks_duplicate.py @@ -1,14 +1,11 @@ -def task1(): - pass - +from .main import fluxqueue -task1.task_name = "task-1" # type: ignore -task1.queue = "default" # type: ignore - -def task2(): +@fluxqueue.task(name="task-1") +def task_1(): pass -task2.task_name = "task-1" # type: ignore -task2.queue = "default" # type: ignore +@fluxqueue.task(name="task-1") +def task_2(): + pass diff --git a/crates/fluxqueue-worker/tests/test_tasks_module.py b/crates/fluxqueue-worker/tests/test_tasks_module.py index 68809f8..22810c9 100644 --- a/crates/fluxqueue-worker/tests/test_tasks_module.py +++ b/crates/fluxqueue-worker/tests/test_tasks_module.py @@ -1,34 +1,25 @@ -def task1(): - pass +from .main import fluxqueue -task1.task_name = "task-1" # type: ignore -task1.queue = "default" # type: ignore +@fluxqueue.task() +def task_1(): + pass -def task2(x: int, y: int): +@fluxqueue.task() +def task_2(x: int, y: int): print(x + y) -task2.task_name = "task-2" # type: ignore -task2.queue = "default" # type: ignore - - +@fluxqueue.task() async def async_task(x: int, y: int): print(x + y) -async_task.task_name = "async-task" # type: ignore -async_task.queue = "default" # type: ignore - - +@fluxqueue.task(queue="high-priority") def high_priority_task(): pass -high_priority_task.task_name = "high-priority-task" # type: ignore -high_priority_task.queue = "high-priority" # type: ignore - - def regular_function(): pass diff --git a/crates/fluxqueue-worker/tests/test_tasks_with_context.py b/crates/fluxqueue-worker/tests/test_tasks_with_context.py new file mode 100644 index 0000000..6dccc28 --- /dev/null +++ b/crates/fluxqueue-worker/tests/test_tasks_with_context.py @@ -0,0 +1,8 @@ +from fluxqueue import Context + +from .main import fluxqueue + + +@fluxqueue.task_with_context() +def sync_func_with_context(ctx: Context): + print(ctx.metadata) diff --git a/python/fluxqueue/__init__.py b/python/fluxqueue/__init__.py index 299142e..de2f3ed 100644 --- a/python/fluxqueue/__init__.py +++ b/python/fluxqueue/__init__.py @@ -1,4 +1,4 @@ -from ._core import __version__ as __version__ # type: ignore +from ._core import __version__ as __version__ from .client import FluxQueue as FluxQueue from .context import Context as Context from .models import TaskMetadata as TaskMetadata diff --git a/python/fluxqueue/_core.pyi b/python/fluxqueue/_core.pyi index a1862d3..abc7062 100644 --- a/python/fluxqueue/_core.pyi +++ b/python/fluxqueue/_core.pyi @@ -1,5 +1,7 @@ from typing import Any +__version__: str + class FluxQueueCore: """ High-performance task queue backed by Rust. From d1e607bd1604ded643b5a01d3c300ab96de9968c Mon Sep 17 00:00:00 2001 From: Giorgi Merebashvili Date: Sun, 1 Mar 2026 17:19:37 +0400 Subject: [PATCH 9/9] Add context tests --- crates/fluxqueue-worker/src/worker.rs | 99 +++++++++++++++++++ .../tests/test_tasks_with_context.py | 17 +++- 2 files changed, 115 insertions(+), 1 deletion(-) diff --git a/crates/fluxqueue-worker/src/worker.rs b/crates/fluxqueue-worker/src/worker.rs index 897373d..f527345 100644 --- a/crates/fluxqueue-worker/src/worker.rs +++ b/crates/fluxqueue-worker/src/worker.rs @@ -445,6 +445,105 @@ mod tests { Ok(()) } + #[tokio::test] + async fn test_sync_task_with_context() -> Result<()> { + let module_path_str = get_test_module_path("test_tasks_with_context.py"); + let task_registry = Arc::new(TaskRegistry::new(&module_path_str, "default")?); + let dispatcher_pool = Arc::new(PythonDispatcher::new(task_registry.clone())?); + + let task = task_registry.get_task(Arc::new("sync-func-with-context".to_string())); + assert!(task.is_some()); + + if let Some(task_func) = task { + let task = Task { + id: "test-id".to_string(), + name: "name".to_string(), + args: vec![144], + kwargs: vec![128], + created_at: 0, + retries: 0, + max_retries: 3, + }; + + let result = run_task( + Arc::new("test".to_string()), + dispatcher_pool.clone(), + Arc::new(task), + task_func, + ) + .await; + assert!(!result.is_err()); + } + + Ok(()) + } + + #[tokio::test] + async fn test_async_task_with_context() -> Result<()> { + let module_path_str = get_test_module_path("test_tasks_with_context.py"); + let task_registry = Arc::new(TaskRegistry::new(&module_path_str, "default")?); + let dispatcher_pool = Arc::new(PythonDispatcher::new(task_registry.clone())?); + + let task = task_registry.get_task(Arc::new("async-func-with-context".to_string())); + assert!(task.is_some()); + + if let Some(task_func) = task { + let task = Task { + id: "test-id".to_string(), + name: "name".to_string(), + args: vec![144], + kwargs: vec![128], + created_at: 0, + retries: 0, + max_retries: 3, + }; + + let result = run_task( + Arc::new("test".to_string()), + dispatcher_pool.clone(), + Arc::new(task), + task_func, + ) + .await; + assert!(!result.is_err()); + } + + Ok(()) + } + + #[tokio::test] + async fn test_task_with_custom_context() -> Result<()> { + let module_path_str = get_test_module_path("test_tasks_with_context.py"); + let task_registry = Arc::new(TaskRegistry::new(&module_path_str, "default")?); + let dispatcher_pool = Arc::new(PythonDispatcher::new(task_registry.clone())?); + + let task = task_registry.get_task(Arc::new("test-custom-context".to_string())); + assert!(task.is_some()); + + if let Some(task_func) = task { + let task = Task { + id: "test-id".to_string(), + name: "name".to_string(), + args: vec![144], + kwargs: vec![128], + created_at: 0, + retries: 0, + max_retries: 3, + }; + + let result = run_task( + Arc::new("test".to_string()), + dispatcher_pool.clone(), + Arc::new(task), + task_func, + ) + .await; + assert!(!result.is_err()); + } + + Ok(()) + } + async fn enqueue_tasks(redis_url: &str) -> Result<()> { use deadpool_redis::{Config, Runtime}; use fluxqueue_common::{ diff --git a/crates/fluxqueue-worker/tests/test_tasks_with_context.py b/crates/fluxqueue-worker/tests/test_tasks_with_context.py index 6dccc28..60a2a3e 100644 --- a/crates/fluxqueue-worker/tests/test_tasks_with_context.py +++ b/crates/fluxqueue-worker/tests/test_tasks_with_context.py @@ -5,4 +5,19 @@ @fluxqueue.task_with_context() def sync_func_with_context(ctx: Context): - print(ctx.metadata) + print("sync_func_with_context metadata: ", ctx.metadata) + + +@fluxqueue.task_with_context() +async def async_func_with_context(ctx: Context): + print("async_func_with_context metadata: ", ctx.metadata) + + +class TestContext(Context): + def get_test_connection(self): + return "conn" + + +@fluxqueue.task_with_context() +def test_custom_context(ctx: TestContext): + print("test_custom_context metadata: ", ctx.metadata)